Skip to content

Commit fe27feb

Browse files
authored
[Dataflow Streaming] Enable state tag encoding v2 (#38705)
* Enable state tag encoding v2 by default for new Dataflow Streaming Engine jobs Also added CHANGES.md entry detailing the change, how to disable it, and job update/downgrade limitations. * Respect UpdateCompatibility for tag encoding v2
1 parent 8b07c7a commit fe27feb

3 files changed

Lines changed: 51 additions & 1 deletion

File tree

CHANGES.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -69,7 +69,7 @@
6969

7070
## New Features / Improvements
7171

72-
* X feature added (Java/Python) ([#X](https://github.com/apache/beam/issues/X)).
72+
* (Java) Enabled state tag encoding v2 by default for new Dataflow Streaming Engine jobs. It can be disabled by passing `--experiments=disable_streaming_engine_state_tag_encoding_v2` or `--updateCompatibilityVersion=2.74.0` pipeline option. Note that the tag encoding version cannot change during a job update. Jobs using tag encoding v2 (enabled by default for new jobs on 2.75.0+) cannot be downgraded to Beam versions prior to 2.73.0, as only versions 2.73.0 and later support tag encoding v2. ([#38705](https://github.com/apache/beam/issues/38705)).
7373

7474
## Breaking Changes
7575

runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunner.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -101,6 +101,7 @@
101101
import org.apache.beam.sdk.options.PipelineOptions;
102102
import org.apache.beam.sdk.options.PipelineOptionsValidator;
103103
import org.apache.beam.sdk.options.SdkHarnessOptions;
104+
import org.apache.beam.sdk.options.StreamingOptions;
104105
import org.apache.beam.sdk.options.ValueProvider.NestedValueProvider;
105106
import org.apache.beam.sdk.runners.AppliedPTransform;
106107
import org.apache.beam.sdk.runners.PTransformOverride;
@@ -1310,6 +1311,11 @@ public DataflowPipelineJob run(Pipeline pipeline) {
13101311
// Experiment marking that the harness supports tag encoding v2
13111312
// Backend will enable tag encoding v2 only if the harness supports it.
13121313
experiments.add("streaming_engine_state_tag_encoding_v2_supported");
1314+
// Experiment requesting tag encoding v2 on new jobs starting with 2.75.0. During job
1315+
// updates old job's tag encoding version is carried over by the backend.
1316+
if (!StreamingOptions.updateCompatibilityVersionLessThan(options, "2.75.0")) {
1317+
experiments.add("enable_streaming_engine_state_tag_encoding_v2");
1318+
}
13131319
options.setExperiments(ImmutableList.copyOf(experiments));
13141320
}
13151321

runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowRunnerTest.java

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2923,4 +2923,48 @@ public void processElement(
29232923
PAssert.that(output).containsInAnyOrder("value:UPDATE_BEFORE");
29242924
pipeline.run();
29252925
}
2926+
2927+
@Test
2928+
public void testStreamingStateTagEncodingV2PreCompatibility() throws Exception {
2929+
DataflowPipelineOptions options = buildPipelineOptions();
2930+
options.as(StreamingOptions.class).setStreaming(true);
2931+
options.as(StreamingOptions.class).setUpdateCompatibilityVersion("2.74.0");
2932+
Pipeline p = Pipeline.create(options);
2933+
2934+
p.run();
2935+
2936+
List<String> experiments = options.getExperiments();
2937+
assertNotNull(experiments);
2938+
assertTrue(experiments.contains("streaming_engine_state_tag_encoding_v2_supported"));
2939+
assertFalse(experiments.contains("enable_streaming_engine_state_tag_encoding_v2"));
2940+
}
2941+
2942+
@Test
2943+
public void testStreamingStateTagEncodingV2PostCompatibility() throws Exception {
2944+
DataflowPipelineOptions options = buildPipelineOptions();
2945+
options.as(StreamingOptions.class).setStreaming(true);
2946+
options.as(StreamingOptions.class).setUpdateCompatibilityVersion("2.75.0");
2947+
Pipeline p = Pipeline.create(options);
2948+
2949+
p.run();
2950+
2951+
List<String> experiments = options.getExperiments();
2952+
assertNotNull(experiments);
2953+
assertTrue(experiments.contains("streaming_engine_state_tag_encoding_v2_supported"));
2954+
assertTrue(experiments.contains("enable_streaming_engine_state_tag_encoding_v2"));
2955+
}
2956+
2957+
@Test
2958+
public void testStreamingStateTagEncodingV2NoCompatibility() throws Exception {
2959+
DataflowPipelineOptions options = buildPipelineOptions();
2960+
options.as(StreamingOptions.class).setStreaming(true);
2961+
Pipeline p = Pipeline.create(options);
2962+
2963+
p.run();
2964+
2965+
List<String> experiments = options.getExperiments();
2966+
assertNotNull(experiments);
2967+
assertTrue(experiments.contains("streaming_engine_state_tag_encoding_v2_supported"));
2968+
assertTrue(experiments.contains("enable_streaming_engine_state_tag_encoding_v2"));
2969+
}
29262970
}

0 commit comments

Comments
 (0)