Apache Beam JdbcIO弃用withCoder后如何设置Coder?
解决Apache Beam JdbcIO.readAll()无法推断Coder的问题
问题分析
你遇到的错误是因为Beam在JdbcIO.readAll()的扩展阶段无法自动推断输出类型InboundData的Coder,尽管你在PCollection上调用了setCoder,但这个操作时机太晚——JdbcIO内部需要提前知道Coder才能完成扩展逻辑。
推荐解决方案(避免弃用的withCoder)
方案1:给InboundData添加@DefaultCoder注解
直接在数据类上标注默认Coder,Beam会自动识别并使用该Coder,无需手动设置:
import org.apache.beam.sdk.coders.DefaultCoder; import org.apache.beam.sdk.coders.SerializableCoder; @DefaultCoder(SerializableCoder.class) public class InboundData implements Serializable { // 类字段与方法定义 }
方案2:让InboundData实现Serializable接口
如果你的类已经实现了Serializable,Beam会自动将SerializableCoder作为该类型的默认Coder,无需额外注解:
public class InboundData implements Serializable { // 类字段与方法定义 }
方案3:全局注册Coder
如果需要在整个Pipeline中为InboundData统一指定Coder,可以通过CoderRegistry注册:
Pipeline pipeline = Pipeline.create(options); pipeline.getCoderRegistry().registerCoder(InboundData.class, SerializableCoder.of(InboundData.class));
该方式适合多个转换都用到InboundData的场景。
为什么之前的setCoder无效?
PCollection.setCoder()是在JdbcIO转换完成后执行的,但JdbcIO.readAll()在自身的expand方法执行时就需要确定输出Coder,此时setCoder还未生效,因此会抛出无法推断的错误。优先使用注解或Serializable实现的方式,让Beam在扩展阶段就能自动获取Coder。
内容的提问来源于stack exchange,提问作者jics
相关产品推荐
相关产品推荐

