diff --git a/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPipelineOptions.java b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPipelineOptions.java index a5a8bb9328b1..e95268ac1308 100644 --- a/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPipelineOptions.java +++ b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/KafkaStreamsPipelineOptions.java @@ -49,7 +49,12 @@ public interface KafkaStreamsPipelineOptions extends PortablePipelineOptions { void setMaxBundleSize(int maxBundleSize); - @Description("Soft cap on bundle wall-clock duration in milliseconds.") + @Description( + "Intended cap on how long a bundle may stay open, in milliseconds. NOT APPLIED YET: closing a" + + " bundle from a wall-clock punctuator made a pipeline with two chained GroupByKeys" + + " across several partitions emit its groups repeatedly against a real broker, so only" + + " the element-count bound is enforced for now. See" + + " https://github.com/apache/beam/issues/18479.") @Default.Integer(1000) int getMaxBundleTimeMs(); diff --git a/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageProcessor.java b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageProcessor.java index 63fefe3e15c4..ef376606114b 100644 --- a/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageProcessor.java +++ b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageProcessor.java @@ -65,6 +65,18 @@ * across the upstream transform's partitions actually advances. Until every partition has reported, * the watermark is held and nothing is forwarded — but data is still processed in the meantime. * + *

A bundle is also bounded in size, by {@code --maxBundleSize}, and closed once that many + * elements have been fed to it. Without the bound a bundle stays open until the next watermark, + * which on a stream that produces steadily lets it grow without limit. The bound is checked as + * elements arrive. A time bound ({@code --maxBundleTimeMs}) is not applied yet — see the option's + * own documentation. + * + *

Closing a bundle asks Kafka Streams to commit, so the elements a bundle consumed and the + * records it produced are committed together and a restart replays either all of the bundle or none + * of it. Note that this aligns commits to bundle boundaries but does not stop Kafka + * Streams from committing on its own interval part-way through a bundle; closing the bundle first + * from a pre-commit hook would be needed to rule that out entirely. + * *

This is the Kafka Streams analogue of Flink's {@code ExecutableStageDoFnOperator} and Spark's * {@code SparkExecutableStageFunction}. State, timers, and side inputs are out of scope for this * first version: the stage is executed with {@link StateRequestHandler#unsupported()} and no timer @@ -107,6 +119,12 @@ class ExecutableStageProcessor private @Nullable StageBundleFactory stageBundleFactory; private @Nullable RemoteBundle currentBundle; + /** Bound on how many elements may be fed to one bundle. */ + private final int maxBundleSize; + + /** Elements fed to the open bundle, for the size bound above. */ + private int elementsInBundle; + /** * @param transformId this stage's own transform id, stamped on the watermarks it emits * @param upstreamTransformIds the transform ids feeding this stage (known from the pipeline @@ -120,13 +138,15 @@ class ExecutableStageProcessor String transformId, Set upstreamTransformIds, MetricsContainerImpl metricsContainer, - Map outputChildByPCollectionId) { + Map outputChildByPCollectionId, + int maxBundleSize) { this.stagePayload = stagePayload; this.jobInfo = jobInfo; this.transformId = transformId; this.watermarkAggregator = new WatermarkAggregator(upstreamTransformIds); this.metricsContainer = metricsContainer; this.outputChildByPCollectionId = ImmutableMap.copyOf(outputChildByPCollectionId); + this.maxBundleSize = maxBundleSize; } /** A harness output element together with the id of the output PCollection it belongs to. */ @@ -185,9 +205,13 @@ public void process(Record> record) { try { ensureBundleOpen(); mainInputReceiver().accept(payload.getData()); + elementsInBundle++; } catch (Exception e) { throw new RuntimeException("Failed to process element through SDK harness", e); } + if (elementsInBundle >= maxBundleSize) { + closeBundleAndFlush(record); + } } private void ensureBundleOpen() throws Exception { @@ -241,6 +265,7 @@ public void onCompleted(ProcessBundleResponse response) { currentBundle = factory.getBundle( outputReceiverFactory, StateRequestHandler.unsupported(), progressHandler); + elementsInBundle = 0; } private FnDataReceiver> mainInputReceiver() { @@ -252,6 +277,18 @@ private FnDataReceiver> mainInputReceiver() { return receiver; } + /** + * Finishes the open bundle, forwards everything it produced, and asks Kafka Streams to commit. + * + *

The commit request is what ties a bundle to a transaction: the elements the bundle consumed + * and the records it produced are then committed together, so a restart either replays the whole + * bundle or none of it. + * + *

The outputs carry the key of the record that closed the bundle. An executable stage is + * unkeyed — it runs stateless, with no state or timers — so the Kafka record key means nothing to + * it and is only being carried along; where the key does matter, downstream sets it, as {@link + * ShuffleByKeyProcessor} does from the Beam key before a GroupByKey. + */ private void closeBundleAndFlush(Record> record) { RemoteBundle bundle = currentBundle; if (bundle == null) { @@ -265,6 +302,7 @@ private void closeBundleAndFlush(Record> record) { throw new RuntimeException("Failed to close SDK harness bundle", e); } finally { currentBundle = null; + elementsInBundle = 0; } ProcessorContext> ctx = checkInitialized(context); // The harness has finished the bundle (close() returned) so no further enqueues happen. @@ -283,6 +321,7 @@ private void closeBundleAndFlush(Record> record) { ctx.forward(outputRecord, childNode); } } + ctx.commit(); } private void forwardWatermark(Record> record, long watermarkMillis) { diff --git a/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageTranslator.java b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageTranslator.java index caefa6534fad..c59e5b919aba 100644 --- a/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageTranslator.java +++ b/runners/kafka-streams/src/main/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageTranslator.java @@ -108,7 +108,8 @@ public void translate( transformId, ImmutableSet.of(parentProcessor), context.getMetricsContainerStepMap().getContainer(transformId), - outputChildByPCollectionId), + outputChildByPCollectionId, + context.getPipelineOptions().getMaxBundleSize()), parentProcessor); if (multiOutput) { diff --git a/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/BundleBoundaryTest.java b/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/BundleBoundaryTest.java new file mode 100644 index 000000000000..c057434f0b2a --- /dev/null +++ b/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/BundleBoundaryTest.java @@ -0,0 +1,118 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.runners.kafka.streams.translation; + +import static org.hamcrest.CoreMatchers.is; +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.greaterThanOrEqualTo; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import org.apache.beam.runners.kafka.streams.KafkaStreamsPipelineOptions; +import org.apache.beam.runners.kafka.streams.KafkaStreamsTestRunner; +import org.apache.beam.sdk.Pipeline; +import org.apache.beam.sdk.transforms.Create; +import org.apache.beam.sdk.transforms.DoFn; +import org.apache.beam.sdk.transforms.ParDo; +import org.junit.Before; +import org.junit.Test; + +/** + * Tests that a bundle is closed once it reaches {@code --maxBundleSize}, rather than staying open + * until the next watermark. + * + *

The bound is observable through the DoFn's own lifecycle: {@code @FinishBundle} runs once per + * bundle the SDK harness processes, so feeding a known number of elements with a known bound tells + * us how many bundles the stage actually opened. + * + *

The elements have to reach the stage as separate records for the bound to see them, since it + * counts what is fed to the stage rather than what the user's code emits inside it. A DoFn that + * fans one element out into many would be fused into the same stage and still be a single input, so + * these pipelines read the elements from a {@link Create} instead — the runner translates that to a + * primitive Read, which forwards one record per element. + */ +public class BundleBoundaryTest { + + private static final int ELEMENTS = 50; + private static final int MAX_BUNDLE_SIZE = 10; + + /** Records how many times {@code @FinishBundle} fired, i.e. how many bundles were processed. */ + private static final List FINISHED_BUNDLES = + Collections.synchronizedList(new ArrayList<>()); + + @Before + public void resetCounters() { + FINISHED_BUNDLES.clear(); + } + + /** Counts the bundles it is asked to process. */ + private static class CountBundlesFn extends DoFn { + @ProcessElement + public void processElement(@Element Integer element, OutputReceiver out) { + out.output(element); + } + + @FinishBundle + public void finishBundle() { + FINISHED_BUNDLES.add("bundle"); + } + } + + private static List elements() { + List elements = new ArrayList<>(); + for (int i = 0; i < ELEMENTS; i++) { + elements.add(i); + } + return elements; + } + + private static Pipeline buildPipeline(KafkaStreamsPipelineOptions options) { + Pipeline pipeline = Pipeline.create(options); + pipeline + .apply("read", Create.of(elements())) + .apply("countBundles", ParDo.of(new CountBundlesFn())); + return pipeline; + } + + private static Pipeline pipelineWithBundleSize(int maxBundleSize) { + KafkaStreamsPipelineOptions options = + KafkaStreamsTestRunner.testOptions().as(KafkaStreamsPipelineOptions.class); + options.setMaxBundleSize(maxBundleSize); + return buildPipeline(options); + } + + @Test + public void aBundleIsClosedOnceItReachesTheSizeBound() { + KafkaStreamsTestRunner.run(pipelineWithBundleSize(MAX_BUNDLE_SIZE)); + + // 50 elements bounded at 10 cannot have gone through in fewer than 5 bundles. Without the + // bound the whole run is one bundle, so this is what tells the two apart. The count is a lower + // bound rather than exact: a watermark arriving mid-bundle also closes one. + assertThat(FINISHED_BUNDLES.size(), is(greaterThanOrEqualTo(ELEMENTS / MAX_BUNDLE_SIZE))); + } + + @Test + public void aBoundLargerThanTheInputLeavesASingleBundle() { + // The control: with a bound nothing reaches, the stage keeps one bundle open until the + // terminal watermark closes it. + KafkaStreamsTestRunner.run(pipelineWithBundleSize(ELEMENTS * 10)); + + assertThat(FINISHED_BUNDLES.size(), is(1)); + } +} diff --git a/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageProcessorWatermarkTest.java b/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageProcessorWatermarkTest.java index 010d98d52e88..290e109796a8 100644 --- a/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageProcessorWatermarkTest.java +++ b/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/ExecutableStageProcessorWatermarkTest.java @@ -62,7 +62,9 @@ private static ExecutableStageProcessor newProcessor() { ImmutableSet.of(UPSTREAM_ID), new MetricsContainerImpl(STAGE_ID), // Single-output: no per-output routing (this test drives the watermark path directly). - ImmutableMap.of()); + ImmutableMap.of(), + // The bundle size bound is irrelevant to the watermark path this test drives. + 1000); } /** A report from the upstream transform's given partition. */ diff --git a/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/MetricsAcrossBundlesTest.java b/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/MetricsAcrossBundlesTest.java new file mode 100644 index 000000000000..a5b690c39ddb --- /dev/null +++ b/runners/kafka-streams/src/test/java/org/apache/beam/runners/kafka/streams/translation/MetricsAcrossBundlesTest.java @@ -0,0 +1,81 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.beam.runners.kafka.streams.translation; + +import static org.hamcrest.CoreMatchers.is; +import static org.hamcrest.MatcherAssert.assertThat; + +import org.apache.beam.runners.kafka.streams.KafkaStreamsPipelineOptions; +import org.apache.beam.runners.kafka.streams.KafkaStreamsTestRunner; +import org.apache.beam.sdk.Pipeline; +import org.apache.beam.sdk.metrics.Counter; +import org.apache.beam.sdk.metrics.MetricNameFilter; +import org.apache.beam.sdk.metrics.MetricQueryResults; +import org.apache.beam.sdk.metrics.MetricResults; +import org.apache.beam.sdk.metrics.Metrics; +import org.apache.beam.sdk.metrics.MetricsFilter; +import org.apache.beam.sdk.transforms.Create; +import org.apache.beam.sdk.transforms.DoFn; +import org.apache.beam.sdk.transforms.ParDo; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables; +import org.junit.Test; + +/** + * Checks that a user counter stays correct when the bundle size bound makes a stage run many small + * bundles instead of one large one. + * + *

The runner folds the metrics the SDK harness reports into the job's step map as each bundle + * completes, and those updates add rather than replace. That is right only if each report covers + * its own bundle, so splitting the same input across more bundles must not change the total. + */ +public class MetricsAcrossBundlesTest { + private static class CountingFn extends DoFn { + private final Counter counter = Metrics.counter("probe", "elements"); + + @ProcessElement + public void processElement(@Element Integer in, OutputReceiver out) { + counter.inc(); + out.output(in); + } + } + + private static long run(int maxBundleSize) { + KafkaStreamsPipelineOptions options = + KafkaStreamsTestRunner.testOptions().as(KafkaStreamsPipelineOptions.class); + options.setMaxBundleSize(maxBundleSize); + Pipeline p = Pipeline.create(options); + p.apply(Create.of(1, 2, 3, 4, 5, 6)).apply(ParDo.of(new CountingFn())); + MetricResults metrics = KafkaStreamsTestRunner.run(p); + MetricQueryResults q = + metrics.queryMetrics( + MetricsFilter.builder() + .addNameFilter(MetricNameFilter.named("probe", "elements")) + .build()); + return Iterables.getOnlyElement(q.getCounters()).getAttempted(); + } + + @Test + public void oneBundle() { + assertThat(run(1000), is(6L)); + } + + @Test + public void manyBundles() { + assertThat(run(1), is(6L)); + } +}