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 =

Reply via email to