51zero / 51zero/eel-sdk

Kafka Avro Schema Option

未关闭
#159 0 条评论 0 个 reaction 已指派 0 人 在 GitHub 查看
enhancement help wanted priority
主要语言
Scala
星标
147
派生
32
PR 合并指标
30 天内没有已合并 PR

描述

It would be great if it would be possible to provide the existing avro schema (from url or local path) option while reading Kafka Avro stream.

Currently running the following code produces the error below

KafkaSource(KafkaSourceConfig("localhost:12345", "consumer"), Set("topic1"), AvroKafkaDeserializer)

Exception in thread "main" java.io.IOException: Not a data file.
at org.apache.avro.file.DataFileStream.initialize(DataFileStream.java:105)
at org.apache.avro.file.DataFileReader.(DataFileReader.java:97)
at io.eels.component.kafka.AvroKafkaDeserializer$.apply(avro.scala:17)
at io.eels.component.kafka.KafkaSource.schema(KafkaSource.scala:36)
at io.eels.FrameSource.schema$lzycompute(Source.scala:29)
at io.eels.FrameSource.schema(Source.scala:29)
at io.eels.plan.SinkPlan$.apply(SinkPlan.scala:17)
at io.eels.Frame$class.to(Frame.scala:362)
at io.eels.FrameSource.to(Source.scala:23)

Thank you.
Roman.

贡献指南

这个仓库没有索引到贡献指南

调研方向

Look at AvroKafkaDeserializer in avro.scala around line 17, where it tries to read the Avro data. The issue is that the deserializer expects a data file but receives raw Avro records. Investigate how to accept a schema URL or local path, possibly using Avro's Schema.Parser. Check KafkaSource.scala line 36 to see how the schema is obtained. A test with a local Kafka topic and Avro schema would verify the fix.

由索引模型根据 Issue 内容生成。

评估

技术栈
kafka, scala
领域
data-engineering, stream-processing
Issue 类型
功能
难度
3/5
预计耗时
1-2 天
活跃度
停滞
描述清晰度
基本清楚
新手友好度
45/100

把新 issue 发到你的邮箱

精选适合新手参与的 GitHub issue 摘要。