hudi-agent commented on code in PR #19416:
URL: https://github.com/apache/hudi/pull/19416#discussion_r3686657520
##########
hudi-utilities/src/test/java/org/apache/hudi/utilities/callback/TestKafkaCallbackProvider.java:
##########
@@ -72,20 +84,58 @@ public static void cleanupClass() throws IOException {
@Test
public void testCallbackMessage() {
- testUtils.createTopic(testTopicName, 2);
-
- HoodieWriteConfig hoodieConfig = createConfigForKafkaCallback();
- HoodieWriteCommitCallback commitCallback =
HoodieCommitCallbackFactory.create(hoodieConfig);
+ int numPartitions = 2;
+ testUtils.createTopic(testTopicName, numPartitions);
List<HoodieWriteStat> stats = generateFakeHoodieWriteStat(1);
- assertDoesNotThrow(() -> commitCallback.call(new
HoodieWriteCommitCallbackMessage(makeNewCommitTime(),
hoodieConfig.getTableName(), hoodieConfig.getBasePath(), stats)));
+ // without a partition config the message is routed by hashing the table
name key
+ HoodieWriteConfig defaultRoutedConfig = createConfigForKafkaCallback(null);
+ HoodieWriteCommitCallback defaultRoutedCallback =
HoodieCommitCallbackFactory.create(defaultRoutedConfig);
+ assertDoesNotThrow(() -> defaultRoutedCallback.call(new
HoodieWriteCommitCallbackMessage(
+ makeNewCommitTime(), defaultRoutedConfig.getTableName(),
defaultRoutedConfig.getBasePath(), stats)));
+
+ // an explicit partition config overrides the key hashing
+ HoodieWriteConfig pinnedConfig = createConfigForKafkaCallback("1");
+ HoodieWriteCommitCallback pinnedCallback =
HoodieCommitCallbackFactory.create(pinnedConfig);
+ assertDoesNotThrow(() -> pinnedCallback.call(new
HoodieWriteCommitCallbackMessage(
+ makeNewCommitTime(), pinnedConfig.getTableName(),
pinnedConfig.getBasePath(), stats)));
+
+ List<ConsumerRecord<String, String>> consumed =
consumeCallbackMessages(numPartitions, 2);
+ // hashing the table name key routes to partition 0, so partition 1 can
only come from the config
+ assertEquals(Arrays.asList(0, 1),
+
consumed.stream().map(ConsumerRecord::partition).sorted().collect(Collectors.toList()));
+ }
+
+ private List<ConsumerRecord<String, String>> consumeCallbackMessages(int
numPartitions, int expectedCount) {
+ Properties consumerProps = new Properties();
+ consumerProps.setProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,
testUtils.brokerAddress());
+ consumerProps.setProperty(ConsumerConfig.GROUP_ID_CONFIG,
"test-kafka-callback-" + UUID.randomUUID());
+ consumerProps.setProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName());
+ consumerProps.setProperty(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName());
+
+ List<ConsumerRecord<String, String>> records = new ArrayList<>();
+ try (KafkaConsumer<String, String> consumer = new
KafkaConsumer<>(consumerProps)) {
+ List<TopicPartition> partitions = IntStream.range(0, numPartitions)
+ .mapToObj(partition -> new TopicPartition(testTopicName, partition))
+ .collect(Collectors.toList());
+ consumer.assign(partitions);
+ consumer.seekToBeginning(partitions);
Review Comment:
🤖 nit: could you express this as `TimeUnit.SECONDS.toMillis(60)` (or a named
constant like `POLL_TIMEOUT_MS`) so the unit is self-documenting?
<sub><i>⚠️ AI-generated; verify before applying. React 👍/👎 to flag
quality.</i></sub>
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]