[ https://issues.apache.org/jira/browse/BEAM-10529?focusedWorklogId=769911&page=com.atlassian.jira.plugin.system.issuetabpanels:worklog-tabpanel#worklog-769911 ]
ASF GitHub Bot logged work on BEAM-10529: ----------------------------------------- Author: ASF GitHub Bot Created on: 12/May/22 21:41 Start Date: 12/May/22 21:41 Worklog Time Spent: 10m Work Description: aaltay commented on PR #17319: URL: https://github.com/apache/beam/pull/17319#issuecomment-1125446706 Could this be merged? Issue Time Tracking ------------------- Worklog Id: (was: 769911) Time Spent: 26h 20m (was: 26h 10m) > Kafka XLang fails for ?empty? key/values > ---------------------------------------- > > Key: BEAM-10529 > URL: https://issues.apache.org/jira/browse/BEAM-10529 > Project: Beam > Issue Type: Bug > Components: cross-language, io-java-kafka > Reporter: Luke Cwik > Assignee: John Casey > Priority: P1 > Fix For: 2.38.0 > > Time Spent: 26h 20m > Remaining Estimate: 0h > > It looks like the Javadoc for ByteArrayDeserializer and StringDeserializer > can return null[1, 2] and we aren't using > NullableCoder.of(ByteArrayCoder.of()) in the expansion[3]. Note that KafkaIO > does this correctly in its regular coder inference logic[4]. > 1: > [https://kafka.apache.org/21/javadoc/org/apache/kafka/common/serialization/ByteArrayDeserializer.html#deserialize-java.lang.String-byte:A-|https://kafka.apache.org/21/javadoc/org/apache/kafka/common/serialization/ByteArrayDeserializer.html#deserialize-java.lang.String-byte:A-2:] > [2:|https://kafka.apache.org/21/javadoc/org/apache/kafka/common/serialization/ByteArrayDeserializer.html#deserialize-java.lang.String-byte:A-2:] > > [https://kafka.apache.org/21/javadoc/org/apache/kafka/common/serialization/StringDeserializer.html#deserialize-java.lang.String-byte:A-] > 3: > [https://github.com/apache/beam/blob/af2d6b0379d64b522ecb769d88e9e7e7b8900208/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java#L478] > 4: > [https://github.com/apache/beam/blob/af2d6b0379d64b522ecb769d88e9e7e7b8900208/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/LocalDeserializerProvider.java#L85] -- This message was sent by Atlassian Jira (v8.20.7#820007)