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);
}