This is an automated email from the ASF dual-hosted git repository.
stankiewicz 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 e3ff1977173 Tag the Kafka write error output with the schema it
actually emits (#39760)
e3ff1977173 is described below
commit e3ff1977173a3b037469331e76a6d5b04d652d97
Author: ZIHAN DAI <[email protected]>
AuthorDate: Fri Sep 25 17:25:11 2026 +1000
Tag the Kafka write error output with the schema it actually emits (#39760)
* Tag the Kafka write error output with the schema it actually emits
ErrorCounterFn emits ErrorHandling.errorRecord(errorSchema, ...), where
errorSchema is already ErrorHandling.errorSchema(inputSchema). The error
PCollection was tagged with the wrapper applied a second time, declaring
a shape no emitted element can match.
Same defect as apache/beam#39759 in TFRecordWriteSchemaTransformProvider.
The existing tests apply ErrorCounterFn directly and tag the output
themselves, so none of them reach the transform's expand().
Signed-off-by: Zihan Dai <[email protected]>
* Add CHANGES.md entry for the Kafka write error schema fix
Signed-off-by: Zihan Dai <[email protected]>
---------
Signed-off-by: Zihan Dai <[email protected]>
Co-authored-by: Zihan Dai <[email protected]>
---
CHANGES.md | 1 +
.../kafka/KafkaWriteSchemaTransformProvider.java | 3 +--
.../KafkaWriteSchemaTransformProviderTest.java | 24 ++++++++++++++++++++++
3 files changed, 26 insertions(+), 2 deletions(-)
diff --git a/CHANGES.md b/CHANGES.md
index cc4456bbaae..982af563be2 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -81,6 +81,7 @@
## Bugfixes
* (Go) Fixed a data race on the Prism runner's artifact cache map in
JobServices ([#32656](https://github.com/apache/beam/issues/32656)).
+* (Java) Fixed the declared schema of the error output of the Kafka write
SchemaTransform, which wrapped the error schema a second time and did not match
the rows it emits ([#39760](https://github.com/apache/beam/issues/39760)).
* Fixed X (Java/Python) ([#X](https://github.com/apache/beam/issues/X)).
## Security Fixes
diff --git
a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaWriteSchemaTransformProvider.java
b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaWriteSchemaTransformProvider.java
index d8d0fe478e3..e59159ba5b8 100644
---
a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaWriteSchemaTransformProvider.java
+++
b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaWriteSchemaTransformProvider.java
@@ -294,8 +294,7 @@ public class KafkaWriteSchemaTransformProvider
}
// TODO: include output from KafkaIO Write once updated from PDone
- PCollection<Row> errorOutput =
-
outputTuple.get(ERROR_TAG).setRowSchema(ErrorHandling.errorSchema(errorSchema));
+ PCollection<Row> errorOutput =
outputTuple.get(ERROR_TAG).setRowSchema(errorSchema);
return PCollectionRowTuple.of(
handleErrors ? configuration.getErrorHandling().getOutput() :
"errors", errorOutput);
}
diff --git
a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaWriteSchemaTransformProviderTest.java
b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaWriteSchemaTransformProviderTest.java
index cc5bd82a5a8..ef53ff0bb83 100644
---
a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaWriteSchemaTransformProviderTest.java
+++
b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaWriteSchemaTransformProviderTest.java
@@ -50,6 +50,7 @@ import org.apache.beam.sdk.transforms.ParDo;
import org.apache.beam.sdk.transforms.SerializableFunction;
import org.apache.beam.sdk.values.KV;
import org.apache.beam.sdk.values.PCollection;
+import org.apache.beam.sdk.values.PCollectionRowTuple;
import org.apache.beam.sdk.values.PCollectionTuple;
import org.apache.beam.sdk.values.Row;
import org.apache.beam.sdk.values.TupleTag;
@@ -280,6 +281,29 @@ public class KafkaWriteSchemaTransformProviderTest {
}
}
+ @Test
+ public void testErrorOutputCarriesTheSchemaErrorCounterFnEmits() {
+ // The output schema is fixed while the graph is built, so this needs no
runner.
+ p.enableAbandonedNodeEnforcement(false);
+
+ Schema inputSchema = Schema.builder().addByteArrayField("bytes").build();
+ KafkaWriteSchemaTransformProvider.KafkaWriteSchemaTransformConfiguration
configuration =
+
KafkaWriteSchemaTransformProvider.KafkaWriteSchemaTransformConfiguration.builder()
+ .setFormat("RAW")
+ .setTopic("test-topic")
+ .setBootstrapServers("host:9092")
+
.setErrorHandling(ErrorHandling.builder().setOutput("errors").build())
+ .build();
+
+ PCollectionRowTuple output =
+ PCollectionRowTuple.of("input", p.apply(Create.empty(inputSchema)))
+ .apply(new
KafkaWriteSchemaTransformProvider().from(configuration));
+
+ // ErrorCounterFn emits ErrorHandling.errorRecord(errorSchema, ...), where
errorSchema is
+ // already ErrorHandling.errorSchema(inputSchema).
+ assertEquals(ErrorHandling.errorSchema(inputSchema),
output.get("errors").getSchema());
+ }
+
@Test
public void testKafkaWriteSchemaTransformConfigurationSchema() throws
NoSuchSchemaException {
Schema schema =