[
https://issues.apache.org/jira/browse/BEAM-7216?focusedWorklogId=237787&page=com.atlassian.jira.plugin.system.issuetabpanels:worklog-tabpanel#worklog-237787
]
ASF GitHub Bot logged work on BEAM-7216:
----------------------------------------
Author: ASF GitHub Bot
Created on: 06/May/19 13:56
Start Date: 06/May/19 13:56
Worklog Time Spent: 10m
Work Description: autodidacticon commented on pull request #8503:
[BEAM-7216] reinstate checks for Kafka client methods
URL: https://github.com/apache/beam/pull/8503
BEAM-7216
In beam 2.9.0, KafkaRecordCoder was used for both producer/consumer records
in KafkaIO, in version 2.10.0, ProducerRecordCoder was introduced but it
appears that in the following code checks are not made to ensure kafka client
compatibility:
https://github.com/apache/beam/blob/master/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/ProducerRecordCoder.java#L137
Specifically the method call to `headers` will fail for kafka clients <
0.11. Elsewhere in this class there are checks on ConsumerSpEL and it is
proposed that they should be reused in the line referenced.
------------------------
Thank you for your contribution! Follow this checklist to help us
incorporate your contribution quickly and easily:
- [ ] [**Choose
reviewer(s)**](https://beam.apache.org/contribute/#make-your-change) and
mention them in a comment (`R: @username`).
- [ ] Format the pull request title like `[BEAM-XXX] Fixes bug in
ApproximateQuantiles`, where you replace `BEAM-XXX` with the appropriate JIRA
issue, if applicable. This will automatically link the pull request to the
issue.
- [ ] If this contribution is large, please file an Apache [Individual
Contributor License Agreement](https://www.apache.org/licenses/icla.pdf).
Post-Commit Tests Status (on master branch)
------------------------------------------------------------------------------------------------
Lang | SDK | Apex | Dataflow | Flink | Gearpump | Samza | Spark
--- | --- | --- | --- | --- | --- | --- | ---
Go | [](https://builds.apache.org/job/beam_PostCommit_Go/lastCompletedBuild/)
| --- | --- | --- | --- | --- | ---
Java | [](https://builds.apache.org/job/beam_PostCommit_Java/lastCompletedBuild/)
| [](https://builds.apache.org/job/beam_PostCommit_Java_ValidatesRunner_Apex/lastCompletedBuild/)
| [](https://builds.apache.org/job/beam_PostCommit_Java_ValidatesRunner_Dataflow/lastCompletedBuild/)
| [](https://builds.apache.org/job/beam_PostCommit_Java_ValidatesRunner_Flink/lastCompletedBuild/)<br>[](https://builds.apache.org/job/beam_PostCommit_Java_PVR_Flink_Batch/lastCompletedBuild/)<br>[](https://builds.apache.org/job/beam_PostCommit_Java_PVR_Flink_Streaming/lastCompletedBuild/)
| [](https://builds.apache.org/job/beam_PostCommit_Java_ValidatesRunner_Gearpump/lastCompletedBuild/)
| [](https://builds.apache.org/job/beam_PostCommit_Java_ValidatesRunner_Samza/lastCompletedBuild/)
| [](https://builds.apache.org/job/beam_PostCommit_Java_ValidatesRunner_Spark/lastCompletedBuild/)
Python | [](https://builds.apache.org/job/beam_PostCommit_Python_Verify/lastCompletedBuild/)<br>[](https://builds.apache.org/job/beam_PostCommit_Python3_Verify/lastCompletedBuild/)
| --- | [](https://builds.apache.org/job/beam_PostCommit_Py_VR_Dataflow/lastCompletedBuild/)
<br> [](https://builds.apache.org/job/beam_PostCommit_Py_ValCont/lastCompletedBuild/)
| [](https://builds.apache.org/job/beam_PreCommit_Python_PVR_Flink_Cron/lastCompletedBuild/)
| --- | --- | ---
Pre-Commit Tests Status (on master branch)
------------------------------------------------------------------------------------------------
--- |Java | Python | Go | Website
--- | --- | --- | --- | ---
Non-portable | [](https://builds.apache.org/job/beam_PreCommit_Java_Cron/lastCompletedBuild/)
| [](https://builds.apache.org/job/beam_PreCommit_Python_Cron/lastCompletedBuild/)
| [](https://builds.apache.org/job/beam_PreCommit_Go_Cron/lastCompletedBuild/)
| [](https://builds.apache.org/job/beam_PreCommit_Website_Cron/lastCompletedBuild/)
Portable | --- | [](https://builds.apache.org/job/beam_PreCommit_Portable_Python_Cron/lastCompletedBuild/)
| --- | ---
See
[.test-infra/jenkins/README](https://github.com/apache/beam/blob/master/.test-infra/jenkins/README.md)
for trigger phrase, status and link of all Jenkins jobs.
----------------------------------------------------------------
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
For queries about this service, please contact Infrastructure at:
[email protected]
Issue Time Tracking
-------------------
Worklog Id: (was: 237787)
Time Spent: 10m
Remaining Estimate: 0h
> beam-sdks-java-io-kafka error with kafka brokers < 0.11
> -------------------------------------------------------
>
> Key: BEAM-7216
> URL: https://issues.apache.org/jira/browse/BEAM-7216
> Project: Beam
> Issue Type: Bug
> Components: io-java-kafka
> Affects Versions: 2.10.0, 2.11.0, 2.12.0
> Reporter: Richard Moorhead
> Priority: Minor
> Time Spent: 10m
> Remaining Estimate: 0h
>
> In beam 2.9.0, KafkaRecordCoder was used for both producer/consumer records
> in KafkaIO, in version 2.10.0, ProducerRecordCoder was introduced but it
> appears that in the following code checks are not made to ensure kafka client
> compatibility:
> [https://github.com/apache/beam/blob/master/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/ProducerRecordCoder.java#L137]
> Specifically the method call to `headers` will fail for kafka clients < 0.11.
> Elsewhere in this class there are checks on ConsumerSpEL and it is proposed
> that they should be reused in the line referenced.
--
This message was sent by Atlassian JIRA
(v7.6.3#76005)