You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Apache Flink 1.16中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_

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.07 22:45:41