Skip to content

Commit 79ed7cc

Browse files
authored
Remove legacy processContext usage across and replace it with argument provider (#38366)
Starting with examples for java and kotlin.
1 parent e7cb9f7 commit 79ed7cc

79 files changed

Lines changed: 898 additions & 592 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

examples/java/iceberg/src/main/java/org/apache/beam/examples/iceberg/IcebergBatchWriteExample.java

Lines changed: 8 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,8 @@
2828
import org.apache.beam.sdk.options.Validation;
2929
import org.apache.beam.sdk.schemas.Schema;
3030
import org.apache.beam.sdk.transforms.DoFn;
31+
import org.apache.beam.sdk.transforms.DoFn.Element;
32+
import org.apache.beam.sdk.transforms.DoFn.OutputReceiver;
3133
import org.apache.beam.sdk.transforms.MapElements;
3234
import org.apache.beam.sdk.transforms.PTransform;
3335
import org.apache.beam.sdk.transforms.ParDo;
@@ -74,9 +76,8 @@ private static Row flattenAnalyticsRow(Row row) {
7476

7577
static class ExtractBrowserTransactionsFn extends DoFn<Row, KV<String, Long>> {
7678
@ProcessElement
77-
public void processElement(ProcessContext c) {
78-
Row row = c.element();
79-
c.output(
79+
public void processElement(@Element Row row, OutputReceiver<KV<String, Long>> receiver) {
80+
receiver.output(
8081
KV.of(
8182
Preconditions.checkStateNotNull(row.getString("browser")),
8283
Preconditions.checkStateNotNull(row.getInt64("transactions"))));
@@ -85,13 +86,13 @@ public void processElement(ProcessContext c) {
8586

8687
static class FormatCountsFn extends DoFn<KV<String, Long>, Row> {
8788
@ProcessElement
88-
public void processElement(ProcessContext c) {
89+
public void processElement(@Element KV<String, Long> element, OutputReceiver<Row> receiver) {
8990
Row row =
9091
Row.withSchema(AGGREGATED_SCHEMA)
91-
.withFieldValue("browser", c.element().getKey())
92-
.withFieldValue("transaction_count", c.element().getValue())
92+
.withFieldValue("browser", element.getKey())
93+
.withFieldValue("transaction_count", element.getValue())
9394
.build();
94-
c.output(row);
95+
receiver.output(row);
9596
}
9697
}
9798

examples/java/sql/src/main/java/org/apache/beam/examples/SchemaTransformExample.java

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@
2222
// description: Demonstration of Schema transform usage.
2323
// multifile: false
2424
// default_example: false
25-
// context_line: 60
25+
// context_line: 65
2626
// categories:
2727
// - Schemas
2828
// - Combiners
@@ -42,6 +42,8 @@
4242
import org.apache.beam.sdk.schemas.transforms.Select;
4343
import org.apache.beam.sdk.transforms.Create;
4444
import org.apache.beam.sdk.transforms.DoFn;
45+
import org.apache.beam.sdk.transforms.DoFn.Element;
46+
import org.apache.beam.sdk.transforms.DoFn.OutputReceiver;
4547
import org.apache.beam.sdk.transforms.Max;
4648
import org.apache.beam.sdk.transforms.Min;
4749
import org.apache.beam.sdk.transforms.ParDo;
@@ -101,9 +103,9 @@ public LogOutput(String prefix) {
101103
}
102104

103105
@ProcessElement
104-
public void processElement(ProcessContext c) throws Exception {
105-
LOG.info("{}{}", prefix, c.element());
106-
c.output(c.element());
106+
public void processElement(@Element T element, OutputReceiver<T> receiver) throws Exception {
107+
LOG.info("{}{}", prefix, element);
108+
receiver.output(element);
107109
}
108110
}
109111
}

examples/java/sql/src/main/java/org/apache/beam/examples/SqlTransformExample.java

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@
2222
// description: Demonstration of SQL transform usage.
2323
// multifile: false
2424
// default_example: false
25-
// context_line: 60
25+
// context_line: 62
2626
// categories:
2727
// - Beam SQL
2828
// - Combiners
@@ -41,6 +41,8 @@
4141
import org.apache.beam.sdk.schemas.Schema;
4242
import org.apache.beam.sdk.transforms.Create;
4343
import org.apache.beam.sdk.transforms.DoFn;
44+
import org.apache.beam.sdk.transforms.DoFn.Element;
45+
import org.apache.beam.sdk.transforms.DoFn.OutputReceiver;
4446
import org.apache.beam.sdk.transforms.ParDo;
4547
import org.apache.beam.sdk.values.PCollection;
4648
import org.apache.beam.sdk.values.Row;
@@ -95,9 +97,9 @@ public LogOutput(String prefix) {
9597
}
9698

9799
@ProcessElement
98-
public void processElement(ProcessContext c) throws Exception {
99-
LOG.info("{}{}", prefix, c.element());
100-
c.output(c.element());
100+
public void processElement(@Element T element, OutputReceiver<T> receiver) throws Exception {
101+
LOG.info("{}{}", prefix, element);
102+
receiver.output(element);
101103
}
102104
}
103105
}

examples/java/src/main/java/org/apache/beam/examples/ApproximateQuantilesExample.java

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,8 @@
2424
import org.apache.beam.sdk.transforms.ApproximateQuantiles;
2525
import org.apache.beam.sdk.transforms.Create;
2626
import org.apache.beam.sdk.transforms.DoFn;
27+
import org.apache.beam.sdk.transforms.DoFn.Element;
28+
import org.apache.beam.sdk.transforms.DoFn.OutputReceiver;
2729
import org.apache.beam.sdk.transforms.ParDo;
2830
import org.apache.beam.sdk.values.PCollection;
2931
import org.slf4j.Logger;
@@ -34,7 +36,7 @@
3436
// description: Demonstration of ApproximateQuantiles transform usage.
3537
// multifile: false
3638
// default_example: false
37-
// context_line: 46
39+
// context_line: 48
3840
// categories:
3941
// - Core Transforms
4042
// complexity: BASIC
@@ -70,9 +72,9 @@ public LogOutput(String prefix) {
7072
}
7173

7274
@ProcessElement
73-
public void processElement(ProcessContext c) throws Exception {
74-
LOG.info("{}{}", prefix, c.element());
75-
c.output(c.element());
75+
public void processElement(@Element T element, OutputReceiver<T> receiver) throws Exception {
76+
LOG.info("{}{}", prefix, element);
77+
receiver.output(element);
7678
}
7779
}
7880
}

examples/java/src/main/java/org/apache/beam/examples/CoCombineTransformExample.java

Lines changed: 12 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@
2222
// description: Demonstration of Composed Combine transform usage.
2323
// multifile: false
2424
// default_example: false
25-
// context_line: 143
25+
// context_line: 145
2626
// categories:
2727
// - Schemas
2828
// - Combiners
@@ -46,6 +46,8 @@
4646
import org.apache.beam.sdk.transforms.CombineFns;
4747
import org.apache.beam.sdk.transforms.Create;
4848
import org.apache.beam.sdk.transforms.DoFn;
49+
import org.apache.beam.sdk.transforms.DoFn.Element;
50+
import org.apache.beam.sdk.transforms.DoFn.OutputReceiver;
4951
import org.apache.beam.sdk.transforms.Max;
5052
import org.apache.beam.sdk.transforms.Min;
5153
import org.apache.beam.sdk.transforms.ParDo;
@@ -185,13 +187,16 @@ public Long apply(Long input) {
185187
new DoFn<
186188
KV<Long, CombineFns.CoCombineResult>, KV<Long, Iterable<KV<String, Long>>>>() {
187189
@ProcessElement
188-
public void processElement(ProcessContext c) throws Exception {
189-
CombineFns.CoCombineResult e = c.element().getValue();
190+
public void processElement(
191+
@Element KV<Long, CombineFns.CoCombineResult> element,
192+
OutputReceiver<KV<Long, Iterable<KV<String, Long>>>> receiver)
193+
throws Exception {
194+
CombineFns.CoCombineResult e = element.getValue();
190195
ArrayList<KV<String, Long>> o = new ArrayList<KV<String, Long>>();
191196
o.add(KV.of(minTag.getId(), e.get(minTag)));
192197
o.add(KV.of(maxTag.getId(), e.get(maxTag)));
193198
o.add(KV.of(sumTag.getId(), e.get(sumTag)));
194-
c.output(KV.of(c.element().getKey(), o));
199+
receiver.output(KV.of(element.getKey(), o));
195200
}
196201
}));
197202

@@ -210,9 +215,9 @@ public LogOutput(String prefix) {
210215
}
211216

212217
@ProcessElement
213-
public void processElement(ProcessContext c) throws Exception {
214-
LOG.info("{}{}", prefix, c.element());
215-
c.output(c.element());
218+
public void processElement(@Element T element, OutputReceiver<T> receiver) throws Exception {
219+
LOG.info("{}{}", prefix, element);
220+
receiver.output(element);
216221
}
217222
}
218223
}

examples/java/src/main/java/org/apache/beam/examples/CoGroupByKeyExample.java

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,8 @@
2222
import org.apache.beam.sdk.options.PipelineOptionsFactory;
2323
import org.apache.beam.sdk.transforms.Create;
2424
import org.apache.beam.sdk.transforms.DoFn;
25+
import org.apache.beam.sdk.transforms.DoFn.Element;
26+
import org.apache.beam.sdk.transforms.DoFn.OutputReceiver;
2527
import org.apache.beam.sdk.transforms.ParDo;
2628
import org.apache.beam.sdk.transforms.join.CoGbkResult;
2729
import org.apache.beam.sdk.transforms.join.CoGroupByKey;
@@ -37,7 +39,7 @@
3739
// description: Demonstration of CoGroupByKey transform usage.
3840
// multifile: false
3941
// default_example: false
40-
// context_line: 54
42+
// context_line: 56
4143
// categories:
4244
// - Core Transforms
4345
// - Joins
@@ -84,9 +86,9 @@ public LogOutput(String prefix) {
8486
}
8587

8688
@ProcessElement
87-
public void processElement(ProcessContext c) throws Exception {
88-
LOG.info("{}{}", prefix, c.element());
89-
c.output(c.element());
89+
public void processElement(@Element T element, OutputReceiver<T> receiver) throws Exception {
90+
LOG.info("{}{}", prefix, element);
91+
receiver.output(element);
9092
}
9193
}
9294
}

examples/java/src/main/java/org/apache/beam/examples/CombineExample.java

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,8 @@
2323
import org.apache.beam.sdk.transforms.Combine;
2424
import org.apache.beam.sdk.transforms.Create;
2525
import org.apache.beam.sdk.transforms.DoFn;
26+
import org.apache.beam.sdk.transforms.DoFn.Element;
27+
import org.apache.beam.sdk.transforms.DoFn.OutputReceiver;
2628
import org.apache.beam.sdk.transforms.ParDo;
2729
import org.apache.beam.sdk.transforms.Sum;
2830
import org.apache.beam.sdk.values.PCollection;
@@ -34,7 +36,7 @@
3436
// description: Demonstration of Combine transform usage.
3537
// multifile: false
3638
// default_example: false
37-
// context_line: 47
39+
// context_line: 49
3840
// categories:
3941
// - Core Transforms
4042
// - Combiners
@@ -68,9 +70,9 @@ public LogOutput(String prefix) {
6870
}
6971

7072
@ProcessElement
71-
public void processElement(ProcessContext c) throws Exception {
72-
LOG.info("{}{}", prefix, c.element());
73-
c.output(c.element());
73+
public void processElement(@Element T element, OutputReceiver<T> receiver) throws Exception {
74+
LOG.info("{}{}", prefix, element);
75+
receiver.output(element);
7476
}
7577
}
7678
}

examples/java/src/main/java/org/apache/beam/examples/CountExample.java

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,8 @@
2323
import org.apache.beam.sdk.transforms.Count;
2424
import org.apache.beam.sdk.transforms.Create;
2525
import org.apache.beam.sdk.transforms.DoFn;
26+
import org.apache.beam.sdk.transforms.DoFn.Element;
27+
import org.apache.beam.sdk.transforms.DoFn.OutputReceiver;
2628
import org.apache.beam.sdk.transforms.ParDo;
2729
import org.apache.beam.sdk.values.PCollection;
2830
import org.slf4j.Logger;
@@ -33,7 +35,7 @@
3335
// description: Demonstration of Count transform usage.
3436
// multifile: false
3537
// default_example: false
36-
// context_line: 45
38+
// context_line: 47
3739
// categories:
3840
// - Core Transforms
3941
// complexity: BASIC
@@ -63,9 +65,9 @@ public LogOutput(String prefix) {
6365
}
6466

6567
@ProcessElement
66-
public void processElement(ProcessContext c) throws Exception {
67-
LOG.info("{}{}", prefix, c.element());
68-
c.output(c.element());
68+
public void processElement(@Element T element, OutputReceiver<T> receiver) throws Exception {
69+
LOG.info("{}{}", prefix, element);
70+
receiver.output(element);
6971
}
7072
}
7173
}

examples/java/src/main/java/org/apache/beam/examples/CountPerKeyExample.java

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,8 @@
2323
import org.apache.beam.sdk.transforms.Count;
2424
import org.apache.beam.sdk.transforms.Create;
2525
import org.apache.beam.sdk.transforms.DoFn;
26+
import org.apache.beam.sdk.transforms.DoFn.Element;
27+
import org.apache.beam.sdk.transforms.DoFn.OutputReceiver;
2628
import org.apache.beam.sdk.transforms.ParDo;
2729
import org.apache.beam.sdk.values.KV;
2830
import org.apache.beam.sdk.values.PCollection;
@@ -34,7 +36,7 @@
3436
// description: Demonstration of Count.perKey transform usage.
3537
// multifile: false
3638
// default_example: false
37-
// context_line: 47
39+
// context_line: 49
3840
// categories:
3941
// - Core Transforms
4042
// complexity: BASIC
@@ -67,9 +69,9 @@ public LogOutput(String prefix) {
6769
}
6870

6971
@ProcessElement
70-
public void processElement(ProcessContext c) throws Exception {
71-
LOG.info("{}{}", prefix, c.element());
72-
c.output(c.element());
72+
public void processElement(@Element T element, OutputReceiver<T> receiver) throws Exception {
73+
LOG.info("{}{}", prefix, element);
74+
receiver.output(element);
7375
}
7476
}
7577
}

examples/java/src/main/java/org/apache/beam/examples/CreateExample.java

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,8 @@
2727
import org.apache.beam.sdk.options.PipelineOptionsFactory;
2828
import org.apache.beam.sdk.transforms.Create;
2929
import org.apache.beam.sdk.transforms.DoFn;
30+
import org.apache.beam.sdk.transforms.DoFn.Element;
31+
import org.apache.beam.sdk.transforms.DoFn.OutputReceiver;
3032
import org.apache.beam.sdk.transforms.ParDo;
3133
import org.apache.beam.sdk.values.KV;
3234
import org.apache.beam.sdk.values.PCollection;
@@ -38,7 +40,7 @@
3840
// description: Demonstration of Create transform usage.
3941
// multifile: false
4042
// default_example: false
41-
// context_line: 51
43+
// context_line: 53
4244
// categories:
4345
// - Core Transforms
4446
// complexity: BASIC
@@ -79,9 +81,9 @@ public LogOutput(String prefix) {
7981
}
8082

8183
@ProcessElement
82-
public void processElement(ProcessContext c) throws Exception {
83-
LOG.info("{}{}", prefix, c.element());
84-
c.output(c.element());
84+
public void processElement(@Element T element, OutputReceiver<T> receiver) throws Exception {
85+
LOG.info("{}{}", prefix, element);
86+
receiver.output(element);
8587
}
8688
}
8789
}

0 commit comments

Comments
 (0)