This is an automated email from the ASF dual-hosted git repository.
Jackie-Jiang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pinot.git
The following commit(s) were added to refs/heads/master by this push:
new 50ab2681b36 Fix flaky waitForMinionTaskCompletion by waiting for
subtask states to become terminal (#19265)
50ab2681b36 is described below
commit 50ab2681b36f9f85f7a5639d37bcfc8ed8a93c58
Author: Xiaotian (Jackie) Jiang <[email protected]>
AuthorDate: Sat Aug 15 19:29:04 2026 -0700
Fix flaky waitForMinionTaskCompletion by waiting for subtask states to
become terminal (#19265)
---
.../tests/BaseClusterIntegrationTest.java | 154 ++++++++++++---------
1 file changed, 87 insertions(+), 67 deletions(-)
diff --git
a/pinot-integration-test-base/src/test/java/org/apache/pinot/integration/tests/BaseClusterIntegrationTest.java
b/pinot-integration-test-base/src/test/java/org/apache/pinot/integration/tests/BaseClusterIntegrationTest.java
index e2574a9b709..ba8218a01ce 100644
---
a/pinot-integration-test-base/src/test/java/org/apache/pinot/integration/tests/BaseClusterIntegrationTest.java
+++
b/pinot-integration-test-base/src/test/java/org/apache/pinot/integration/tests/BaseClusterIntegrationTest.java
@@ -64,6 +64,7 @@ import
org.apache.pinot.common.restlet.resources.ValidDocIdsType;
import org.apache.pinot.common.utils.LLCSegmentName;
import org.apache.pinot.common.utils.TarCompressionUtils;
import org.apache.pinot.common.utils.config.TagNameUtils;
+import
org.apache.pinot.controller.helix.core.minion.PinotHelixTaskResourceManager;
import org.apache.pinot.plugin.inputformat.csv.CSVMessageDecoder;
import org.apache.pinot.plugin.stream.kafka.KafkaStreamConfigProperties;
import org.apache.pinot.plugin.stream.kafka30.server.EmbeddedKafkaCluster;
@@ -96,7 +97,10 @@ import org.apache.pinot.spi.utils.builder.TableNameBuilder;
import org.apache.pinot.tools.utils.KafkaStarterUtils;
import org.apache.pinot.util.TestUtils;
import org.intellij.lang.annotations.Language;
-import org.testng.Assert;
+
+import static org.testng.Assert.assertNotNull;
+import static org.testng.Assert.assertNull;
+import static org.testng.Assert.assertTrue;
/// Shared implementation details of the cluster integration tests.
@@ -317,7 +321,7 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
protected Schema createSchema(String schemaFileName)
throws IOException {
InputStream inputStream =
getClass().getClassLoader().getResourceAsStream(schemaFileName);
- Assert.assertNotNull(inputStream);
+ assertNotNull(inputStream);
return Schema.fromInputStream(inputStream);
}
@@ -329,22 +333,20 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
protected TableConfig createTableConfig(String tableConfigFileName)
throws IOException {
URL configPathUrl =
getClass().getClassLoader().getResource(tableConfigFileName);
- Assert.assertNotNull(configPathUrl);
+ assertNotNull(configPathUrl);
return createTableConfig(new File(configPathUrl.getFile()));
}
protected TableConfig createTableConfig(File tableConfigFile)
throws IOException {
InputStream inputStream = new FileInputStream(tableConfigFile);
- Assert.assertNotNull(inputStream);
+ assertNotNull(inputStream);
return JsonUtils.inputStreamToObject(inputStream, TableConfig.class);
}
/// Creates a new OFFLINE table config.
protected TableConfig createOfflineTableConfig() {
- // @formatter:off
- return new TableConfigBuilder(TableType.OFFLINE)
- .setTableName(getTableName())
+ return new
TableConfigBuilder(TableType.OFFLINE).setTableName(getTableName())
.setTimeColumnName(getTimeColumnName())
.setSortedColumn(getSortedColumn())
.setInvertedIndexColumns(getInvertedIndexColumns())
@@ -364,7 +366,6 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
.setSegmentPartitionConfig(getSegmentPartitionConfig())
.setOptimizeNoDictStatsCollection(true)
.build();
- // @formatter:on
}
/// Returns the OFFLINE table config in the cluster.
@@ -381,15 +382,13 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
String streamType = "kafka";
streamConfigMap.put(StreamConfigProperties.STREAM_TYPE, streamType);
streamConfigMap.put(KafkaStreamConfigProperties.constructStreamProperty(
- KafkaStreamConfigProperties.LowLevelConsumer.KAFKA_BROKER_LIST),
- getKafkaBrokerList());
+ KafkaStreamConfigProperties.LowLevelConsumer.KAFKA_BROKER_LIST),
getKafkaBrokerList());
if (useKafkaTransaction()) {
streamConfigMap.put(KafkaStreamConfigProperties.constructStreamProperty(
KafkaStreamConfigProperties.LowLevelConsumer.KAFKA_ISOLATION_LEVEL),
KafkaStreamConfigProperties.LowLevelConsumer.KAFKA_ISOLATION_LEVEL_READ_COMMITTED);
// Ensure the consumer can fetch complete transactional batches plus
commit markers.
- streamConfigMap.put(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG,
- Integer.toString(10 * 1024 * 1024));
+ streamConfigMap.put(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG,
Integer.toString(10 * 1024 * 1024));
}
streamConfigMap.put(StreamConfigProperties.constructStreamProperty(streamType,
StreamConfigProperties.STREAM_CONSUMER_FACTORY_CLASS),
getStreamConsumerFactoryClassName());
@@ -436,14 +435,12 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
physicalTableConfigMap.put(offlineTableName, new PhysicalTableConfig());
physicalTableConfigMap.put(realtimeTableName, new PhysicalTableConfig());
- return new LogicalTableConfigBuilder()
- .setTableName(getLogicalTableName())
+ return new LogicalTableConfigBuilder().setTableName(getLogicalTableName())
.setBrokerTenant(getBrokerTenant())
.setRefOfflineTableName(offlineTableName)
.setRefRealtimeTableName(realtimeTableName)
.setPhysicalTableConfigMap(physicalTableConfigMap)
- .setTimeBoundaryConfig(
- new TimeBoundaryConfig("min", Map.of("includedTables",
physicalTableConfigMap.keySet())))
+ .setTimeBoundaryConfig(new TimeBoundaryConfig("min",
Map.of("includedTables", physicalTableConfigMap.keySet())))
.build();
}
@@ -457,8 +454,7 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
// TODO - Use this method to create table config for all table types to
avoid redundant code
protected TableConfigBuilder getTableConfigBuilder(TableType tableType) {
- return new TableConfigBuilder(tableType)
- .setTableName(getTableName())
+ return new TableConfigBuilder(tableType).setTableName(getTableName())
.setTimeColumnName(getTimeColumnName())
.setSortedColumn(getSortedColumn())
.setInvertedIndexColumns(getInvertedIndexColumns())
@@ -494,15 +490,24 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
upsertConfig.setDeleteRecordColumn(deleteColumn);
return new
TableConfigBuilder(TableType.REALTIME).setTableName(getTableName())
-
.setTimeColumnName(getTimeColumnName()).setFieldConfigList(getFieldConfigs()).setNumReplicas(getNumReplicas())
-
.setSegmentVersion(getSegmentVersion()).setLoadMode(getLoadMode()).setTaskConfig(getTaskConfig())
-
.setBrokerTenant(getBrokerTenant()).setServerTenant(getServerTenant()).setIngestionConfig(getIngestionConfig())
-
.setStreamConfigs(getStreamConfigs()).setNullHandlingEnabled(getNullHandlingEnabled()).setRoutingConfig(
+ .setTimeColumnName(getTimeColumnName())
+ .setFieldConfigList(getFieldConfigs())
+ .setNumReplicas(getNumReplicas())
+ .setSegmentVersion(getSegmentVersion())
+ .setLoadMode(getLoadMode())
+ .setTaskConfig(getTaskConfig())
+ .setBrokerTenant(getBrokerTenant())
+ .setServerTenant(getServerTenant())
+ .setIngestionConfig(getIngestionConfig())
+ .setStreamConfigs(getStreamConfigs())
+ .setNullHandlingEnabled(getNullHandlingEnabled())
+ .setRoutingConfig(
new RoutingConfig(null, null,
RoutingConfig.STRICT_REPLICA_GROUP_INSTANCE_SELECTOR_TYPE, false))
.setSegmentPartitionConfig(new
SegmentPartitionConfig(columnPartitionConfigMap))
.setReplicaGroupStrategyConfig(new
ReplicaGroupStrategyConfig(primaryKeyColumn, 1))
.setOptimizeNoDictStatsCollection(true)
- .setUpsertConfig(upsertConfig).build();
+ .setUpsertConfig(upsertConfig)
+ .build();
}
protected Map<String, String> getCSVDecoderProperties(@Nullable String
delimiter,
@@ -543,17 +548,25 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
kafkaTopicName);
streamConfigsMap.putAll(streamDecoderProperties);
- return new
TableConfigBuilder(TableType.REALTIME).setTableName(tableName).setTimeColumnName(getTimeColumnName())
-
.setFieldConfigList(getFieldConfigs()).setNumReplicas(getNumReplicas()).setSegmentVersion(getSegmentVersion())
-
.setLoadMode(getLoadMode()).setTaskConfig(getTaskConfig()).setBrokerTenant(getBrokerTenant())
-
.setServerTenant(getServerTenant()).setIngestionConfig(getIngestionConfig()).setStreamConfigs(streamConfigsMap)
+ return new TableConfigBuilder(TableType.REALTIME).setTableName(tableName)
+ .setTimeColumnName(getTimeColumnName())
+ .setFieldConfigList(getFieldConfigs())
+ .setNumReplicas(getNumReplicas())
+ .setSegmentVersion(getSegmentVersion())
+ .setLoadMode(getLoadMode())
+ .setTaskConfig(getTaskConfig())
+ .setBrokerTenant(getBrokerTenant())
+ .setServerTenant(getServerTenant())
+ .setIngestionConfig(getIngestionConfig())
+ .setStreamConfigs(streamConfigsMap)
.setNullHandlingEnabled(UpsertConfig.Mode.PARTIAL.equals(upsertConfig.getMode())
|| getNullHandlingEnabled())
.setRoutingConfig(
new RoutingConfig(null, null,
RoutingConfig.STRICT_REPLICA_GROUP_INSTANCE_SELECTOR_TYPE, false))
.setSegmentPartitionConfig(new
SegmentPartitionConfig(columnPartitionConfigMap))
.setReplicaGroupStrategyConfig(new
ReplicaGroupStrategyConfig(primaryKeyColumn, 1))
.setOptimizeNoDictStatsCollection(true)
- .setUpsertConfig(upsertConfig).build();
+ .setUpsertConfig(upsertConfig)
+ .build();
}
/// Creates a new Dedup enabled table config
@@ -608,11 +621,11 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
return _pinotConnectionV2;
}
if (_pinotConnection == null) {
- JsonAsyncHttpPinotClientTransportFactory factory = new
JsonAsyncHttpPinotClientTransportFactory()
- .withConnectionProperties(getPinotConnectionProperties());
+ JsonAsyncHttpPinotClientTransportFactory factory =
+ new
JsonAsyncHttpPinotClientTransportFactory().withConnectionProperties(getPinotConnectionProperties());
factory.setHeaders(getPinotClientTransportHeaders());
- _pinotConnection = ConnectionFactory.fromZookeeper(getZkUrl() + "/" +
getHelixClusterName(),
- factory.buildTransport());
+ _pinotConnection =
+ ConnectionFactory.fromZookeeper(getZkUrl() + "/" +
getHelixClusterName(), factory.buildTransport());
}
return _pinotConnection;
}
@@ -627,7 +640,7 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
///
/// @return H2 connection
protected Connection getH2Connection() {
- Assert.assertNotNull(_h2Connection, "H2 Connection has not been
initialized");
+ assertNotNull(_h2Connection, "H2 Connection has not been initialized");
return _h2Connection;
}
@@ -635,14 +648,14 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
///
/// @return Query generator.
protected QueryGenerator getQueryGenerator() {
- Assert.assertNotNull(_queryGenerator, "Query Generator has not been
initialized");
+ assertNotNull(_queryGenerator, "Query Generator has not been initialized");
return _queryGenerator;
}
/// Sets up the H2 connection
protected void setUpH2Connection()
throws Exception {
- Assert.assertNull(_h2Connection);
+ assertNull(_h2Connection);
Class.forName("org.h2.Driver");
_h2Connection = DriverManager.getConnection("jdbc:h2:mem:");
}
@@ -656,7 +669,7 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
/// Sets up the query generator using the given Avro files.
protected void setUpQueryGenerator(List<File> avroFiles) {
- Assert.assertNull(_queryGenerator);
+ assertNull(_queryGenerator);
String tableName = getTableName();
_queryGenerator = new QueryGenerator(avroFiles, tableName, tableName);
}
@@ -675,7 +688,7 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
protected List<File> unpackTarData(String tarFileName, File outputDir)
throws Exception {
InputStream inputStream =
getClass().getClassLoader().getResourceAsStream(tarFileName);
- Assert.assertNotNull(inputStream);
+ assertNotNull(inputStream);
return TarCompressionUtils.untar(inputStream, outputDir);
}
@@ -702,7 +715,8 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
}
protected void createAndUploadSegmentFromClasspath(TableConfig tableConfig,
Schema schema, String dataFilePath,
- FileFormat fileFormat, long expectedNoOfDocs, long timeoutMs) throws
Exception {
+ FileFormat fileFormat, long expectedNoOfDocs, long timeoutMs)
+ throws Exception {
URL dataPathUrl = getClass().getClassLoader().getResource(dataFilePath);
assert dataPathUrl != null;
File file = new File(dataPathUrl.getFile());
@@ -714,12 +728,14 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
/// dataFilePath on the classpath
@Deprecated
protected void createAndUploadSegmentFromFile(TableConfig tableConfig,
Schema schema, String dataFilePath,
- FileFormat fileFormat, long expectedNoOfDocs, long timeoutMs) throws
Exception {
+ FileFormat fileFormat, long expectedNoOfDocs, long timeoutMs)
+ throws Exception {
createAndUploadSegmentFromClasspath(tableConfig, schema, dataFilePath,
fileFormat, expectedNoOfDocs, timeoutMs);
}
protected void createAndUploadSegmentFromFile(TableConfig tableConfig,
Schema schema, File file,
- FileFormat fileFormat, long expectedNoOfDocs, long timeoutMs) throws
Exception {
+ FileFormat fileFormat, long expectedNoOfDocs, long timeoutMs)
+ throws Exception {
TestUtils.ensureDirectoriesExistAndEmpty(_segmentDir, _tarDir);
ClusterIntegrationTestUtils.buildSegmentFromFile(file, tableConfig,
schema, "%", _segmentDir, _tarDir, fileFormat);
@@ -790,8 +806,7 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
private void waitForKafkaClusterReady(String brokerList, int brokerCount,
boolean requireTransactions) {
TestUtils.waitForCondition(aVoid -> isKafkaClusterReady(brokerList,
brokerCount), 200L,
- KAFKA_CLUSTER_READY_TIMEOUT_MS,
- "Kafka brokers are not ready");
+ KAFKA_CLUSTER_READY_TIMEOUT_MS, "Kafka brokers are not ready");
if (requireTransactions) {
// Wait for transaction coordinator with longer initial delay and timeout
TestUtils.waitForCondition(aVoid -> canInitTransactions(brokerList),
1000L, KAFKA_CLUSTER_READY_TIMEOUT_MS,
@@ -842,9 +857,8 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
protected void createKafkaTopic(String topic, int numPartitions, int
replicationFactor) {
int effectiveReplicationFactor = Math.max(1, Math.min(replicationFactor,
getNumKafkaBrokers()));
- _kafkaStarters.get(0).createTopic(
- topic,
- KafkaStarterUtils.getTopicCreationProps(numPartitions,
effectiveReplicationFactor));
+ _kafkaStarters.get(0)
+ .createTopic(topic,
KafkaStarterUtils.getTopicCreationProps(numPartitions,
effectiveReplicationFactor));
waitForKafkaTopicReady(topic, numPartitions, effectiveReplicationFactor);
waitForKafkaTopicMetadataReadyForConsumer(topic, numPartitions);
}
@@ -872,15 +886,13 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
}
private void waitForKafkaTopicReady(String topic, int expectedPartitions,
int expectedReplicationFactor) {
- TestUtils.waitForCondition(
- aVoid -> isKafkaTopicReady(topic, expectedPartitions,
expectedReplicationFactor),
- 200L, KAFKA_TOPIC_READY_TIMEOUT_MS, "Kafka topic '" + topic + "' is
not fully ready");
+ TestUtils.waitForCondition(aVoid -> isKafkaTopicReady(topic,
expectedPartitions, expectedReplicationFactor), 200L,
+ KAFKA_TOPIC_READY_TIMEOUT_MS, "Kafka topic '" + topic + "' is not
fully ready");
}
private void waitForKafkaTopicMetadataReadyForConsumer(String topic, int
expectedPartitions) {
TestUtils.waitForCondition(aVoid ->
isKafkaTopicMetadataReadyForConsumer(topic, expectedPartitions), 200L,
- KAFKA_TOPIC_READY_TIMEOUT_MS,
- "Kafka topic '" + topic + "' metadata is not visible to consumers");
+ KAFKA_TOPIC_READY_TIMEOUT_MS, "Kafka topic '" + topic + "' metadata is
not visible to consumers");
}
private boolean isKafkaTopicReady(String topic, int expectedPartitions, int
expectedReplicationFactor) {
@@ -914,8 +926,7 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
adminProps.put(AdminClientConfig.DEFAULT_API_TIMEOUT_MS_CONFIG, "5000");
try (AdminClient adminClient = AdminClient.create(adminProps)) {
TopicDescription topicDescription =
- adminClient.describeTopics(List.of(topic)).allTopicNames().get(5,
TimeUnit.SECONDS)
- .get(topic);
+ adminClient.describeTopics(List.of(topic)).allTopicNames().get(5,
TimeUnit.SECONDS).get(topic);
if (topicDescription.partitions().size() < expectedPartitions) {
return false;
}
@@ -979,19 +990,28 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
}
protected void waitForMinionTaskCompletion(String taskId, long timeout) {
- TestUtils.waitForCondition(aVoid ->
-
_controllerStarter.getHelixTaskResourceManager().getTaskState(taskId) ==
TaskState.COMPLETED,
- timeout, "Failed to complete the task " + taskId);
-
- // Validate that there were > 0 subtasks so that we know the task was
actually run
-
Assert.assertFalse(_controllerStarter.getHelixTaskResourceManager().getSubtaskStates(taskId).isEmpty());
+ // The task state (workflow context) and the subtask states (job context)
live in different Helix znodes and are
+ // not updated atomically, so a subtask can still read as a non-terminal
state (null/INIT before start, RUNNING,
+ // or STOPPED before a resume) right after the task turns COMPLETED. Wait
until the task is COMPLETED and every
+ // subtask reached a terminal state before validating them. A non-empty
subtask map also proves the task was
+ // actually run.
+ PinotHelixTaskResourceManager taskResourceManager =
_controllerStarter.getHelixTaskResourceManager();
+ TestUtils.waitForCondition(aVoid -> {
+ if (taskResourceManager.getTaskState(taskId) != TaskState.COMPLETED) {
+ return false;
+ }
+ Map<String, TaskPartitionState> subtaskStates =
taskResourceManager.getSubtaskStates(taskId);
+ return !subtaskStates.isEmpty() && subtaskStates.values()
+ .stream()
+ .noneMatch(state -> state == null || state ==
TaskPartitionState.INIT || state == TaskPartitionState.RUNNING
+ || state == TaskPartitionState.STOPPED);
+ }, timeout, "Failed to complete the task " + taskId);
// Validate that all subtasks are completed successfully. A task can be
marked completed even if some subtasks
// failed, so we need to check the subtask states.
- Map<String, TaskPartitionState> subTaskStates =
_controllerStarter.getHelixTaskResourceManager()
- .getSubtaskStates(taskId);
- Assert.assertTrue(subTaskStates.values().stream().allMatch(x -> x ==
TaskPartitionState.COMPLETED),
- "Not all subtasks are completed for task " + taskId + " : " +
subTaskStates);
+ Map<String, TaskPartitionState> subtaskStates =
taskResourceManager.getSubtaskStates(taskId);
+ assertTrue(subtaskStates.values().stream().allMatch(state -> state ==
TaskPartitionState.COMPLETED),
+ "Not all subtasks are completed for task " + taskId + " : " +
subtaskStates);
}
protected List<String> getSegments(String tableNameWithType) {
@@ -1049,8 +1069,7 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
continue;
}
Map<String, String> stateMap =
idealState.getInstanceStateMap(segmentName);
- if (stateMap != null
- &&
stateMap.containsValue(CommonConstants.Helix.StateModel.SegmentStateModel.CONSUMING))
{
+ if (stateMap != null &&
stateMap.containsValue(CommonConstants.Helix.StateModel.SegmentStateModel.CONSUMING))
{
consumingPartitions.add(new
LLCSegmentName(segmentName).getPartitionGroupId());
}
}
@@ -1130,8 +1149,9 @@ public abstract class BaseClusterIntegrationTest extends
ClusterTest {
ValidDocIdsType validDocIdsType)
throws Exception {
- String responseString = getOrCreateAdminClient().getTableClient()
- .getValidDocIdsMetadata(tableNameWithType, validDocIdsType.toString());
- return JsonUtils.stringToObject(responseString, new TypeReference<>() { });
+ String responseString =
+
getOrCreateAdminClient().getTableClient().getValidDocIdsMetadata(tableNameWithType,
validDocIdsType.toString());
+ return JsonUtils.stringToObject(responseString, new TypeReference<>() {
+ });
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]