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

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


The following commit(s) were added to refs/heads/master by this push:
     new b4623815617 [Flink] Restore source split logs (#40354)
b4623815617 is described below

commit b4623815617ef1fc3a2fbfd0bd4d5f96c2402a88
Author: Paulius Kuzmickas <[email protected]>
AuthorDate: Wed Sep 30 18:16:06 2026 +0100

    [Flink] Restore source split logs (#40354)
    
    The move of source splitting into FlinkSourceSplitUtils dropped the
    split-count logs from the default and lazy enumerators. Log in the
    shared helpers so every enumerator reports how many splits a source
    produced, and remove the now-duplicate size-based enumerator log.
---
 .../streaming/io/source/FlinkSourceSplitUtils.java   | 20 ++++++++++++++++++--
 .../source/SizeBasedFlinkSourceSplitEnumerator.java  |  5 -----
 2 files changed, 18 insertions(+), 7 deletions(-)

diff --git 
a/runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/FlinkSourceSplitUtils.java
 
b/runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/FlinkSourceSplitUtils.java
index dda99bf53be..26a5dd4d934 100644
--- 
a/runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/FlinkSourceSplitUtils.java
+++ 
b/runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/FlinkSourceSplitUtils.java
@@ -25,11 +25,15 @@ import org.apache.beam.sdk.io.FileBasedSource;
 import org.apache.beam.sdk.io.Source;
 import org.apache.beam.sdk.io.UnboundedSource;
 import org.apache.beam.sdk.options.PipelineOptions;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 /** Shared Beam source sizing and splitting helpers. */
 final class FlinkSourceSplitUtils {
   static final long MEBIBYTE = 1024L * 1024L;
 
+  private static final Logger LOG = 
LoggerFactory.getLogger(FlinkSourceSplitUtils.class);
+
   private FlinkSourceSplitUtils() {}
 
   static <T> long estimateBoundedSourceSize(
@@ -45,13 +49,25 @@ final class FlinkSourceSplitUtils {
       throws Exception {
     long desiredSizeBytes =
         getDesiredSizeBytes(boundedSource, pipelineOptions, numSplits, 
estimatedSizeBytes);
-    return toFlinkSplits(boundedSource.split(desiredSizeBytes, 
pipelineOptions));
+    List<? extends BoundedSource<T>> splits =
+        boundedSource.split(desiredSizeBytes, pipelineOptions);
+    LOG.info(
+        "Split bounded source {} in {} splits (estimated size {} bytes, "
+            + "desired split size {} bytes)",
+        boundedSource,
+        splits.size(),
+        estimatedSizeBytes,
+        desiredSizeBytes);
+    return toFlinkSplits(splits);
   }
 
   static <T> ArrayList<FlinkSourceSplit<T>> splitUnboundedSource(
       UnboundedSource<T, ?> unboundedSource, PipelineOptions pipelineOptions, 
int numSplits)
       throws Exception {
-    return toFlinkSplits(unboundedSource.split(numSplits, pipelineOptions));
+    List<? extends UnboundedSource<T, ?>> splits =
+        unboundedSource.split(numSplits, pipelineOptions);
+    LOG.info("Split source {} to {} splits", unboundedSource, splits);
+    return toFlinkSplits(splits);
   }
 
   static long getDesiredSizeBytes(
diff --git 
a/runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/SizeBasedFlinkSourceSplitEnumerator.java
 
b/runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/SizeBasedFlinkSourceSplitEnumerator.java
index 78b0fb59f2c..ff1a79e3c30 100644
--- 
a/runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/SizeBasedFlinkSourceSplitEnumerator.java
+++ 
b/runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/SizeBasedFlinkSourceSplitEnumerator.java
@@ -131,11 +131,6 @@ final class SizeBasedFlinkSourceSplitEnumerator<T>
     ArrayList<FlinkSourceSplit<T>> splits =
         FlinkSourceSplitUtils.splitBoundedSource(
             boundedSource, pipelineOptions, numSplits, estimatedSizeBytes);
-    LOG.info(
-        "Split bounded source {} into {} splits using {} assignment",
-        boundedSource,
-        splits.size(),
-        selectedMode);
     return new FlinkSourceEnumeratorState<>(selectedMode, splits);
   }
 

Reply via email to