Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run.",
"modification": 1
"modification": 1,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run.",
"modification": 2
"modification": 2,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
3 changes: 2 additions & 1 deletion .github/trigger_files/beam_PostCommit_Java_DataflowV2.json
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
{
"https://github.com/apache/beam/pull/39893": "Fix nullness in BigQueryIO",
"modification": 9,
"https://github.com/apache/beam/pull/35159": "moving WindowedValue and making an interface"
"https://github.com/apache/beam/pull/35159": "moving WindowedValue and making an interface",
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run.",
"modification": 3
"modification": 3,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1 +1,3 @@
{}
{
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
Expand Up @@ -2,5 +2,6 @@
"https://github.com/apache/beam/pull/36138": "Cleanly separating v1 worker and v2 sdk harness container image handling",
"https://github.com/apache/beam/pull/34902": "Introducing OutputBuilder",
"comment": "Modify this file in a trivial way to cause this test suite to run",
"modification": 2
"modification": 2,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run.",
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run",
"modification": 2
"modification": 2,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run.",
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run.",
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run",
"https://github.com/apache/beam/pull/32440": "test new datastream runner for batch"
"modification": 2
"https://github.com/apache/beam/pull/32440": "test new datastream runner for batch",
"modification": 2,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run",
"modification": 2,
"https://github.com/apache/beam/pull/32440": "test new datastream runner for batch"
"https://github.com/apache/beam/pull/32440": "test new datastream runner for batch",
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run",
"modification": 8
"modification": 8,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run.",
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run.",
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run",
"modification": 7
"modification": 7,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run!",
"modification": 5
"modification": 5,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run!",
"modification": 1,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
Expand Up @@ -2,5 +2,6 @@
"https://github.com/apache/beam/pull/34902": "Introducing OutputBuilder",
"comment": "Modify this file in a trivial way to cause this test suite to run",
"https://github.com/apache/beam/pull/31156": "noting that PR #31156 should run this test",
"https://github.com/apache/beam/pull/35159": "moving WindowedValue and making an interface"
"https://github.com/apache/beam/pull/35159": "moving WindowedValue and making an interface",
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run",
"modification": 16
"modification": 16,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run",
"modification": 2
"modification": 2,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run",
"modification": 3
"modification": 3,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run",
"modification": 2
"modification": 2,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run",
"modification": 8
"modification": 8,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run",
"revision": 3
"revision": 3,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
3 changes: 2 additions & 1 deletion .github/trigger_files/beam_PostCommit_XVR_Direct.json
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
{
"modification": 1
"modification": 1,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
3 changes: 2 additions & 1 deletion .github/trigger_files/beam_PostCommit_XVR_Flink.json
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"modification": 3,
"trigger-2026-04-04": "portable_runner expand_sdf opt-in"
"trigger-2026-04-04": "portable_runner expand_sdf opt-in",
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run",
"modification": 2
"modification": 2,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run.",
"modification": 1
}
"comment": "Modify this file in a trivial way to cause this test suite to run.",
"modification": 1,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
{
"modification": 2
}
"modification": 2,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run.",
"modification": 1
"comment": "Modify this file in a trivial way to cause this test suite to run.",
"modification": 1,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
3 changes: 2 additions & 1 deletion .github/trigger_files/beam_PostCommit_XVR_Spark3.json
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
{
"trigger-2026-07-08": "portable_runner expand_sdf opt-in 2"
"trigger-2026-07-08": "portable_runner expand_sdf opt-in 2",
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
3 changes: 2 additions & 1 deletion .github/trigger_files/beam_PostCommit_Yaml_Xlang_Direct.json
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to run",
"revision": 3
"revision": 3,
"https://github.com/apache/beam/pull/39990": "removing dead code from FnApiDoFnRunner"
}
Original file line number Diff line number Diff line change
Expand Up @@ -350,10 +350,6 @@ public final void addRunnerForPTransform(Context context) throws IOException {
case PTransformTranslation.PAR_DO_TRANSFORM_URN:
mainOutputTag = (TupleTag) ParDoTranslation.getMainOutputTag(parDoPayload);
break;
case PTransformTranslation.SPLITTABLE_SPLIT_AND_SIZE_RESTRICTIONS_URN:
mainOutputTag =
new TupleTag(Iterables.getOnlyElement(pTransform.getOutputsMap().keySet()));
break;
default:
throw new IllegalStateException(
String.format("Unknown urn: %s", pTransform.getSpec().getUrn()));
Expand Down Expand Up @@ -444,17 +440,7 @@ public final void addRunnerForPTransform(Context context) throws IOException {
this.doFnInvoker = DoFnInvokers.tryInvokeSetupFor(doFn, pipelineOptions);

this.startBundleArgumentProvider = new StartBundleArgumentProvider();
// Register the appropriate handlers.
switch (pTransform.getSpec().getUrn()) {
case PTransformTranslation.PAR_DO_TRANSFORM_URN:
case PTransformTranslation.SPLITTABLE_PROCESS_SIZED_ELEMENTS_AND_RESTRICTIONS_URN:
addStartFunction.accept(this::startBundle);
break;
case PTransformTranslation.SPLITTABLE_SPLIT_AND_SIZE_RESTRICTIONS_URN:
// startBundle should not be invoked
default:
// no-op
}
addStartFunction.accept(this::startBundle);

String mainInput;
try {
Expand All @@ -474,49 +460,24 @@ public final void addRunnerForPTransform(Context context) throws IOException {
}
break;
case PTransformTranslation.SPLITTABLE_PROCESS_SIZED_ELEMENTS_AND_RESTRICTIONS_URN:
if (doFnSignature.processElement().observesWindow()
|| (doFnSignature.newTracker() != null && doFnSignature.newTracker().observesWindow())
|| (doFnSignature.getSize() != null && doFnSignature.getSize().observesWindow())
|| (doFnSignature.newWatermarkEstimator() != null
&& doFnSignature.newWatermarkEstimator().observesWindow())
|| !sideInputMapping.isEmpty()) {
mainInputConsumer =
new SplittableFnDataReceiver() {
@Override
public void accept(WindowedValue input) throws Exception {
processElementForWindowObservingSizedElementAndRestriction(input);
}
};
this.processContext = new WindowObservingProcessBundleContext();
} else {
mainInputConsumer =
new SplittableFnDataReceiver() {
@Override
public void accept(WindowedValue input) throws Exception {
// TODO(BEAM-10303): Create a variant which is optimized to not observe the
// windows.
processElementForWindowObservingSizedElementAndRestriction(input);
}
};
this.processContext = new WindowObservingProcessBundleContext();
}
// TODO(BEAM-10303): Create a variant which is optimized to not observe the windows when
// neither the DoFn nor its side inputs observe them.
mainInputConsumer =
new SplittableFnDataReceiver() {
@Override
public void accept(WindowedValue input) throws Exception {
processElementForWindowObservingSizedElementAndRestriction(input);
}
};
this.processContext = new WindowObservingProcessBundleContext();
break;
default:
throw new IllegalStateException("Unknown urn: " + pTransform.getSpec().getUrn());
}
addPCollectionConsumer.accept(pTransform.getInputsOrThrow(mainInput), mainInputConsumer);

this.finishBundleArgumentProvider = new FinishBundleArgumentProvider();
switch (pTransform.getSpec().getUrn()) {
case PTransformTranslation.PAR_DO_TRANSFORM_URN:
case PTransformTranslation.SPLITTABLE_PROCESS_SIZED_ELEMENTS_AND_RESTRICTIONS_URN:
addFinishFunction.accept(this::finishBundle);
break;
case PTransformTranslation.SPLITTABLE_SPLIT_AND_SIZE_RESTRICTIONS_URN:
// finishBundle should not be invoked
default:
// no-op
}
addFinishFunction.accept(this::finishBundle);
addTearDownFunction.accept(this::tearDown);

workCompletedShortId =
Expand Down
Loading
Loading