Apache Flink 1.16中Protobuf Any类型反序列化问题求助
Apache Flink 1.16 Protobuf反序列化器处理google.protobuf.Any类型的问题解决
这是Flink官方Protobuf格式在处理导入的标准Protobuf类型时的代码生成bug——Flink错误地将google.protobuf.Any识别为自定义Protobuf类的内部类(com.company.model.proto.Common.Any),导致Janino编译失败。以下是可行的解决方法:
1. 升级Flink版本(推荐)
该bug已在Flink 1.17及后续版本中修复,社区修正了Protobuf标准类型的解析逻辑,能够正确识别google.protobuf.Any这类导入的官方类型。如果业务允许升级版本,这是最彻底的解决方案。
2. 自定义Protobuf反序列化器(临时兼容方案)
若暂时无法升级,可自定义反序列化器手动处理MessageWrapper,规避Flink自动生成错误代码的问题:
实现自定义反序列化器
public class CustomMessageWrapperDeserializationSchema implements DeserializationSchema<RowData> { private transient ProtoToRowConverter converter; private final MessageWrapper defaultInstance = MessageWrapper.getDefaultInstance(); @Override public void open(InitializationContext context) { // 构建与表结构匹配的类型信息 TypeInformation<RowData> rowType = Types.ROW( Types.STRING(), // kk_trade_id Types.ROW(Types.BYTES()) // message.value ); PbRowType pbRowType = PbRowType.from(rowType); converter = new ProtoToRowConverter(defaultInstance.getDescriptorForType(), pbRowType); } @Override public RowData deserialize(byte[] message) throws IOException { MessageWrapper wrapper = MessageWrapper.parseFrom(message); // 手动处理Any字段,转换为RowData结构 GenericRowData rowData = new GenericRowData(2); // 注意:如果kk_trade_id来自Kafka Key,需根据实际逻辑获取对应值 // rowData.setField(0, ...); if (wrapper.hasMessage()) { GenericRowData messageRow = new GenericRowData(1); messageRow.setField(0, wrapper.getMessage().getValue().toByteArray()); rowData.setField(1, messageRow); } return rowData; } @Override public boolean isEndOfStream(RowData nextElement) { return false; } }
修改表定义使用自定义反序列化器
CREATE TABLE trade_records ( kk_trade_id STRING, message ROW<`value` BYTES>, PRIMARY KEY (kk_trade_id) NOT ENFORCED ) WITH ( 'connector' = 'upsert-kafka', 'value.deserializer' = 'com.company.model.proto.CustomMessageWrapperDeserializationSchema', 'key.format' = 'raw', 'key.fields-prefix' = 'kk_', 'value.fields-include' = 'EXCEPT_KEY' -- 其他配置项 );
3. 临时调整Protobuf定义(应急方案,不推荐)
通过嵌套消息包装Any类型,改变数据结构以规避bug,但会影响上游生产者和下游业务适配,仅作为应急使用:
修改Protobuf定义
syntax = "proto2"; option java_package = "com.company.model.proto"; import "google/protobuf/any.proto"; message WrappedAny { optional google.protobuf.Any value = 1; } message MessageWrapper { optional WrappedAny message = 1; }
对应表定义调整
CREATE TABLE trade_records ( kk_trade_id STRING, message ROW<`value` ROW<`value` BYTES>>, PRIMARY KEY (kk_trade_id) NOT ENFORCED ) WITH ( 'connector' = 'upsert-kafka', 'value.protobuf.message-class-name' = 'com.company.model.proto.Common$MessageWrapper', 'key.format' = 'raw', 'key.fields-prefix' = 'kk_', 'value.fields-include' = 'EXCEPT_KEY', 'value.format' = 'protobuf' -- 其他配置项 );
内容的提问来源于stack exchange,提问作者dv_
相关产品推荐
相关产品推荐

