This is an automated email from the ASF dual-hosted git repository.

voonhous pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git


The following commit(s) were added to refs/heads/master by this push:
     new 91f3869416a8 refactor(flink): simplify streaming write pipeline 
construction (#19579)
91f3869416a8 is described below

commit 91f3869416a85ff81a398818e5445aa2cc1e8f90
Author: Danny Chan <[email protected]>
AuthorDate: Thu Aug 13 11:22:11 2026 +0800

    refactor(flink): simplify streaming write pipeline construction (#19579)
---
 .../java/org/apache/hudi/sink/utils/Pipelines.java | 251 +++++++++++++--------
 .../org/apache/hudi/sink/utils/TestPipelines.java  |  18 ++
 2 files changed, 177 insertions(+), 92 deletions(-)

diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java
index ea2bf868c19d..7bf09e60234b 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java
@@ -411,11 +411,11 @@ public class Pipelines {
     final boolean globalIndex = conf.get(FlinkOptions.INDEX_GLOBAL_ENABLED);
     if (overwrite || OptionsResolver.isBucketIndexType(conf)) {
       return rowDataToHoodieRecord(conf, rowType, dataStream);
-    } else if (bounded && !globalIndex && 
OptionsResolver.isPartitionedTable(conf)) {
+    }
+    if (bounded && !globalIndex && OptionsResolver.isPartitionedTable(conf)) {
       return boundedBootstrap(conf, rowType, dataStream);
-    } else {
-      return streamBootstrap(conf, rowType, dataStream, bounded);
     }
+    return streamBootstrap(conf, rowType, dataStream, bounded);
   }
 
   private static DataStream<HoodieFlinkInternalRow> streamBootstrap(
@@ -501,86 +501,153 @@ public class Pipelines {
                                                       RowType rowType,
                                                       
DataStream<HoodieFlinkInternalRow> dataStream) {
     if (OptionsResolver.isBucketIndexType(conf)) {
-      HoodieIndex.BucketIndexEngineType bucketIndexEngineType = 
OptionsResolver.getBucketEngineType(conf);
-      switch (bucketIndexEngineType) {
-        case SIMPLE:
-          // [HUDI-9036] BucketIndexPartitioner is also used in bulk insert 
mode,
-          // keep use of HoodieKey here in partitionCustom for now
-          Partitioner<HoodieKey> partitioner = 
BucketIndexPartitionerFactory.create(conf);
-          SingleOutputStreamOperator<RowData> bucketWriteStream = dataStream
-              .partitionCustom(
-                  partitioner,
-                  record -> new HoodieKey(record.getRecordKey(), 
record.getPartitionPath()))
-              .transform(
-                  opName("bucket_write", conf),
-                  TypeInformation.of(RowData.class),
-                  BucketStreamWriteOperator.getFactory(conf, rowType))
-              .uid(opUID("bucket_write", conf))
-              .setParallelism(conf.get(FlinkOptions.WRITE_TASKS));
-          declareManagedMemoryIfNecessary(conf, bucketWriteStream, () -> 
OptionsResolver.getWriteBufferSizeInBytes(conf));
-          return bucketWriteStream;
-        case CONSISTENT_HASHING:
-          if (OptionsResolver.isInsertOverwrite(conf)) {
-            // TODO support insert overwrite for consistent bucket index
-            throw new HoodieException("Consistent hashing bucket index does 
not work with insert overwrite using FLINK engine. Use simple bucket index or 
Spark engine.");
-          }
-          SingleOutputStreamOperator<RowData> consistentBucketWriteStream = 
dataStream
-              .transform(
-                  opName("consistent_bucket_assigner", conf),
-                  new HoodieFlinkInternalRowTypeInfo(rowType),
-                  new ProcessOperator<>(new 
ConsistentBucketAssignFunction(conf)))
-              .uid(opUID("consistent_bucket_assigner", conf))
-              .setParallelism(conf.get(FlinkOptions.BUCKET_ASSIGN_TASKS))
-              .keyBy(HoodieFlinkInternalRow::getFileId)
-              .transform(
-                  opName("consistent_bucket_write", conf),
-                  TypeInformation.of(RowData.class),
-                  BucketStreamWriteOperator.getFactory(conf, rowType))
-              .uid(opUID("consistent_bucket_write", conf))
-              .setParallelism(conf.get(FlinkOptions.WRITE_TASKS));
-          declareManagedMemoryIfNecessary(conf, consistentBucketWriteStream, 
() -> OptionsResolver.getWriteBufferSizeInBytes(conf));
-          return consistentBucketWriteStream;
-        default:
-          throw new HoodieNotSupportedException("Unknown bucket index engine 
type: " + bucketIndexEngineType);
-      }
-    } else {
-      String writeOperatorUid = opUID("stream_write", conf);
-      // uuid is used to generate operator id for the write operator, then the 
bucket assign operator can send
-      // operator event to the coordinator of the write operator based on the 
operator id.
-      // @see org.apache.flink.runtime.jobgraph.tasks.TaskOperatorEventGateway.
-      DataStream<HoodieFlinkInternalRow> bucketAssignStream = 
createBucketAssignStream(dataStream, conf, rowType, writeOperatorUid);
-      boolean isStreamingIndexWriteEnabled = 
OptionsResolver.isStreamingIndexWriteEnabled(conf);
-      SingleOutputStreamOperator<RowData> writeDatastream =
-          bucketAssignStream
-              // shuffle by fileId(bucket id)
-              .keyBy(HoodieFlinkInternalRow::getFileId)
-              .transform(
-                  opName("stream_write", conf),
-                  isStreamingIndexWriteEnabled ? 
InternalTypeInfo.of(IndexRowUtils.INDEX_ROW_TYPE) : 
TypeInformation.of(RowData.class),
-                  StreamWriteOperator.getFactory(conf, rowType))
-              .uid(writeOperatorUid)
-              .setParallelism(conf.get(FlinkOptions.WRITE_TASKS));
-      declareManagedMemoryIfNecessary(conf, writeDatastream, () -> 
OptionsResolver.getWriteBufferSizeInBytes(conf));
-      if (isStreamingIndexWriteEnabled) {
-        // index writing pipeline
-        SingleOutputStreamOperator<RowData> indexWriteDatastream = 
writeDatastream
-            .partitionCustom(
-                OptionsResolver.isRecordLevelIndex(conf)
-                    ? new RecordIndexPartitioner(conf)
-                    : new GlobalRecordIndexPartitioner(conf),
-                IndexRowUtils::getHoodieKey)
+      return bucketStreamWrite(conf, rowType, dataStream);
+    }
+
+    String writeOperatorUid = opUID("stream_write", conf);
+    // uuid is used to generate operator id for the write operator, then the 
bucket assign operator can send
+    // operator event to the coordinator of the write operator based on the 
operator id.
+    // @see org.apache.flink.runtime.jobgraph.tasks.TaskOperatorEventGateway.
+    DataStream<HoodieFlinkInternalRow> bucketAssignStream =
+        createBucketAssignStream(dataStream, conf, rowType, writeOperatorUid);
+    boolean isStreamingIndexWriteEnabled = 
OptionsResolver.isStreamingIndexWriteEnabled(conf);
+    SingleOutputStreamOperator<RowData> writeDataStream = bucketAssignStream
+        // shuffle by fileId(bucket id)
+        .keyBy(HoodieFlinkInternalRow::getFileId)
+        .transform(
+            opName("stream_write", conf),
+            isStreamingIndexWriteEnabled
+                ? InternalTypeInfo.of(IndexRowUtils.INDEX_ROW_TYPE)
+                : TypeInformation.of(RowData.class),
+            StreamWriteOperator.getFactory(conf, rowType))
+        .uid(writeOperatorUid)
+        .setParallelism(conf.get(FlinkOptions.WRITE_TASKS));
+    declareManagedMemoryIfNecessary(
+        conf, writeDataStream, () -> 
OptionsResolver.getWriteBufferSizeInBytes(conf));
+
+    return isStreamingIndexWriteEnabled
+        ? addIndexWrite(conf, writeDataStream, writeOperatorUid)
+        : writeDataStream;
+  }
+
+  /**
+   * The bucket index streaming write pipeline.
+   *
+   * <p>For the simple bucket index, the input dataset shuffles directly by 
bucket ID. For the
+   * consistent hashing bucket index, the bucket assigner first assigns a file 
group, then the
+   * dataset shuffles by file ID before passing around to the write function. 
The pipelines look
+   * like the following:
+   *
+   * <pre>
+   * Simple bucket index:
+   *      | input1 | ===\     /=== | task1 |
+   *                   shuffle(by bucket ID)
+   *      | input2 | ===/     \=== | task2 |
+   *
+   * Consistent hashing bucket index:
+   *      | input1 | === | bucket assigner1 | ===\     /=== | task1 |
+   *                                            shuffle(by file ID)
+   *      | input2 | === | bucket assigner2 | ===/     \=== | task2 |
+   *
+   *      Note: a file group must be handled by one write task to avoid write 
conflict.
+   * </pre>
+   *
+   * @param conf       The configuration
+   * @param rowType    The logical row type of the input records
+   * @param dataStream The input data stream
+   * @return the bucket write data stream
+   */
+  private static DataStream<RowData> bucketStreamWrite(
+      Configuration conf,
+      RowType rowType,
+      DataStream<HoodieFlinkInternalRow> dataStream) {
+    HoodieIndex.BucketIndexEngineType bucketIndexEngineType =
+        OptionsResolver.getBucketEngineType(conf);
+    DataStream<HoodieFlinkInternalRow> bucketAssignedStream;
+    String writeOperatorName;
+    switch (bucketIndexEngineType) {
+      case SIMPLE:
+        // [HUDI-9036] BucketIndexPartitioner is also used in bulk insert mode,
+        // keep use of HoodieKey here in partitionCustom for now
+        Partitioner<HoodieKey> partitioner = 
BucketIndexPartitionerFactory.create(conf);
+        bucketAssignedStream = dataStream.partitionCustom(
+            partitioner,
+            record -> new HoodieKey(record.getRecordKey(), 
record.getPartitionPath()));
+        writeOperatorName = "bucket_write";
+        break;
+      case CONSISTENT_HASHING:
+        if (OptionsResolver.isInsertOverwrite(conf)) {
+          // TODO support insert overwrite for consistent bucket index
+          throw new HoodieException("Consistent hashing bucket index does not 
work with insert overwrite using FLINK engine. Use simple bucket index or Spark 
engine.");
+        }
+        bucketAssignedStream = dataStream
             .transform(
-                opName("index_write", conf),
-                TypeInformation.of(RowData.class),
-                new IndexWriteOperator(conf, 
OperatorIDGenerator.fromUid(writeOperatorUid)))
-            .uid(opUID("index_write", conf))
-            .setParallelism(conf.get(FlinkOptions.INDEX_WRITE_TASKS));
-        declareManagedMemoryIfNecessary(conf, indexWriteDatastream, () -> 
conf.get(FlinkOptions.INDEX_RLI_WRITE_BUFFER_SIZE) * 1024L * 1024L);
-        return indexWriteDatastream;
-      } else {
-        return writeDatastream;
-      }
+                opName("consistent_bucket_assigner", conf),
+                new HoodieFlinkInternalRowTypeInfo(rowType),
+                new ProcessOperator<>(new 
ConsistentBucketAssignFunction(conf)))
+            .uid(opUID("consistent_bucket_assigner", conf))
+            .setParallelism(conf.get(FlinkOptions.BUCKET_ASSIGN_TASKS))
+            .keyBy(HoodieFlinkInternalRow::getFileId);
+        writeOperatorName = "consistent_bucket_write";
+        break;
+      default:
+        throw new HoodieNotSupportedException(
+            "Unknown bucket index engine type: " + bucketIndexEngineType);
     }
+
+    SingleOutputStreamOperator<RowData> bucketWriteStream = 
bucketAssignedStream
+        .transform(
+            opName(writeOperatorName, conf),
+            TypeInformation.of(RowData.class),
+            BucketStreamWriteOperator.getFactory(conf, rowType))
+        .uid(opUID(writeOperatorName, conf))
+        .setParallelism(conf.get(FlinkOptions.WRITE_TASKS));
+    declareManagedMemoryIfNecessary(
+        conf, bucketWriteStream, () -> 
OptionsResolver.getWriteBufferSizeInBytes(conf));
+    return bucketWriteStream;
+  }
+
+  /**
+   * The streaming index write pipeline.
+   *
+   * <p>Index rows emitted by the data write operator are routed to the task 
responsible for their
+   * record-index file group. The data write operator UID identifies its 
coordinator so the index
+   * writer can participate in the same commit. The whole pipeline looks like 
the following:
+   *
+   * <pre>
+   *      | data write1 | ===\     /=== | index task1 |
+   *                        shuffle(by index file group)
+   *      | data write2 | ===/     \=== | index task2 |
+   *
+   *      Note: a record-index file group must be handled by one index write 
task.
+   * </pre>
+   *
+   * @param conf             The configuration
+   * @param writeDataStream  The index rows emitted by the data write operator
+   * @param writeOperatorUid The UID of the upstream data write operator
+   * @return the index write data stream
+   */
+  private static DataStream<RowData> addIndexWrite(
+      Configuration conf,
+      DataStream<RowData> writeDataStream,
+      String writeOperatorUid) {
+    SingleOutputStreamOperator<RowData> indexWriteDataStream = writeDataStream
+        .partitionCustom(
+            OptionsResolver.isRecordLevelIndex(conf)
+                ? new RecordIndexPartitioner(conf)
+                : new GlobalRecordIndexPartitioner(conf),
+            IndexRowUtils::getHoodieKey)
+        .transform(
+            opName("index_write", conf),
+            TypeInformation.of(RowData.class),
+            new IndexWriteOperator(conf, 
OperatorIDGenerator.fromUid(writeOperatorUid)))
+        .uid(opUID("index_write", conf))
+        .setParallelism(conf.get(FlinkOptions.INDEX_WRITE_TASKS));
+    declareManagedMemoryIfNecessary(
+        conf,
+        indexWriteDataStream,
+        () -> conf.get(FlinkOptions.INDEX_RLI_WRITE_BUFFER_SIZE) * 1024L * 
1024L);
+    return indexWriteDataStream;
   }
 
   /**
@@ -603,7 +670,8 @@ public class Pipelines {
               new MiniBatchBucketAssignOperator(new 
MinibatchBucketAssignFunction(conf), 
OperatorIDGenerator.fromUid(writeOperatorUid)))
           .uid(opUID(assignerOperatorName, conf))
           .setParallelism(conf.get(FlinkOptions.BUCKET_ASSIGN_TASKS));
-    } else if (OptionsResolver.isRecordLevelIndex(conf)) {
+    }
+    if (OptionsResolver.isRecordLevelIndex(conf)) {
       return inputStream
           .keyBy(HoodieFlinkInternalRow::getRecordKey)
           .transform(
@@ -612,17 +680,16 @@ public class Pipelines {
               new DynamicBucketAssignOperator(new 
DynamicBucketAssignFunction(conf), 
OperatorIDGenerator.fromUid(writeOperatorUid)))
           .uid(opUID(assignerOperatorName, conf))
           .setParallelism(conf.get(FlinkOptions.BUCKET_ASSIGN_TASKS));
-    } else {
-      return inputStream
-          // Key-by record key, to avoid multiple subtasks write to a bucket 
at the same time
-          .keyBy(HoodieFlinkInternalRow::getRecordKey)
-          .transform(
-              assignerOperatorName,
-              new HoodieFlinkInternalRowTypeInfo(rowType),
-              new KeyedProcessOperator<>(new BucketAssignFunction(conf)))
-          .uid(opUID(assignerOperatorName, conf))
-          .setParallelism(conf.get(FlinkOptions.BUCKET_ASSIGN_TASKS));
     }
+    return inputStream
+        // Key-by record key, to avoid multiple subtasks write to a bucket at 
the same time
+        .keyBy(HoodieFlinkInternalRow::getRecordKey)
+        .transform(
+            assignerOperatorName,
+            new HoodieFlinkInternalRowTypeInfo(rowType),
+            new KeyedProcessOperator<>(new BucketAssignFunction(conf)))
+        .uid(opUID(assignerOperatorName, conf))
+        .setParallelism(conf.get(FlinkOptions.BUCKET_ASSIGN_TASKS));
   }
 
   /**
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/TestPipelines.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/TestPipelines.java
index 32a9067f90ce..09340cd270ec 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/TestPipelines.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/TestPipelines.java
@@ -131,6 +131,24 @@ public class TestPipelines {
     
assertTrue(pipeline.getTransformation().getName().startsWith("bucket_write:"));
   }
 
+  @Test
+  void testConsistentBucketWritePipeline() {
+    Configuration conf = defaultConf();
+    conf.set(FlinkOptions.INDEX_TYPE, HoodieIndex.IndexType.BUCKET.name());
+    conf.set(FlinkOptions.BUCKET_INDEX_ENGINE_TYPE,
+        HoodieIndex.BucketIndexEngineType.CONSISTENT_HASHING.name());
+    conf.set(FlinkOptions.BUCKET_ASSIGN_TASKS, 2);
+    conf.set(FlinkOptions.WRITE_TASKS, 4);
+
+    DataStream<RowData> pipeline = Pipelines.hoodieStreamWrite(
+        conf, TestConfigurations.ROW_TYPE, hoodieRowInput());
+
+    assertEquals(4, pipeline.getParallelism());
+    
assertTrue(pipeline.getTransformation().getName().startsWith("consistent_bucket_write:"));
+    assertTrue(transformationNames(pipeline).stream()
+        .anyMatch(name -> name.startsWith("consistent_bucket_assigner:")));
+  }
+
   @Test
   void testBulkInsertAndAppendValidation() {
     Configuration recordIndexConf = defaultConf();

Reply via email to