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]

Reply via email to