[
https://issues.apache.org/jira/browse/BEAM-10529?focusedWorklogId=752418&page=com.atlassian.jira.plugin.system.issuetabpanels:worklog-tabpanel#worklog-752418
]
ASF GitHub Bot logged work on BEAM-10529:
-----------------------------------------
Author: ASF GitHub Bot
Created on: 04/Apr/22 18:03
Start Date: 04/Apr/22 18:03
Worklog Time Spent: 10m
Work Description: TheNeuralBit commented on PR #16923:
URL: https://github.com/apache/beam/pull/16923#issuecomment-1087854834
@jrmccluskey This was my fault, I assumed the Go failure was the one
@lostluck was referencing
[here](https://github.com/apache/beam/pull/16923#issuecomment-1084793890).
Issue Time Tracking
-------------------
Worklog Id: (was: 752418)
Time Spent: 22.5h (was: 22h 20m)
> 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
> Time Spent: 22.5h
> 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.1#820001)
