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

Apache Flink MongoDB连接器读取数据类型不匹配问题求助

问题

编写Apache Flink MongoDB连接器实现JSON数据读写,Sink已成功插入文档,但Source读取时抛出org.bson.BsonInvalidOperationException,提示Value expected to be of type INT64 is of unexpected type OBJECT_ID。排查发现document.getFirstKey()获取到了MongoDB自动生成的_id字段(OBJECT_ID类型),而非存储的整数字段,需要解决如何正确获取文档中的整数值。

Sink代码

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
        
List<Tuple2<String, Integer>> data = new ArrayList<>();
data.add(new Tuple2<>("Hello", 1));
data.add(new Tuple2<>("Hi", 2));
data.add(new Tuple2<>("Hey", 3));
        
DataStream<Tuple2<String, Integer>> stream = env.fromCollection(data);

MongoSink<Tuple2<String, Integer>> sink = MongoSink.<Tuple2<String, Integer>>builder()
            .setUri("mongodb://127.0.0.1:27017")
            .setDatabase("test_db")
            .setCollection("test_coll")
            .setBatchSize(1000)
            .setBatchIntervalMs(1000)
            .setMaxRetries(3)
            .setDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE)
            .setSerializationSchema(
                    (input, context) 
                        -> {
                            Document doc = new Document(input.f0, input.f1);
                            return new InsertOneModel<>(BsonDocument.parse(doc.toJson()));
                        })
            .build();

stream.sinkTo(sink);

插入的文档示例

{
  "_id": {
    "$oid": "65f67f3b9779060fd2390d0e"
  },
  "Hello": 1
}

Source代码

MongoSource<Tuple2<String,Integer>> source = MongoSource.<Tuple2<String,Integer>>builder()
            .setUri("mongodb://127.0.0.1:27017")
            .setDatabase("test_db")
            .setCollection("test_coll")
            .setDeserializationSchema(new MongoDeserializationSchema<Tuple2<String, Integer>>() {
                @Override
                public Tuple2<String, Integer> deserialize(BsonDocument document) {
                    String key = document.getFirstKey();
                    Integer value = document.getInt64(key).intValue();  // 该行抛出异常

                    return new Tuple2<String, Integer>(key, value);
                }

                @Override
                public TypeInformation<Tuple2<String, Integer>> getProducedType() {
                    return Types.TUPLE(Types.STRING, Types.INT);
                }
            })
            .build();

DataStream<Tuple2<String, Integer>> ds = env.fromSource(source, WatermarkStrategy.noWatermarks(), "MongoDB-Source");
ds.print();

异常信息

Caused by: org.bson.BsonInvalidOperationException: Value expected to be of type INT64 is of unexpected type OBJECT_ID
        at org.bson.BsonValue.throwIfInvalidType(BsonValue.java:419)
        at org.bson.BsonValue.asInt64(BsonValue.java:105)
        at org.bson.BsonDocument.getInt64(BsonDocument.java:203)
        at com.aaa.test.FlinkMongoTest$1.deserialize(FlinkMongoTest.java:63)
        at com.aaa.test.FlinkMongoTest$1.deserialize(FlinkMongoTest.java:1)
        at org.apache.flink.connector.mongodb.source.reader.deserializer.MongoDeserializationSchema.deserialize(MongoDeserializationSchema.java:58)
        at org.apache.flink.connector.mongodb.source.reader.emitter.MongoRecordEmitter.emitRecord(MongoRecordEmitter.java:54) 
        at org.apache.flink.connector.mongodb.source.reader.emitter.MongoRecordEmitter.emitRecord(MongoRecordEmitter.java:34) 
        at org.apache.flink.connector.base.source.reader.SourceReaderBase.pollNext(SourceReaderBase.java:160)
解决方法

方案1:过滤_id字段后提取目标值

在反序列化逻辑中先移除_id字段,再获取剩余的唯一键值对:

@Override
public Tuple2<String, Integer> deserialize(BsonDocument document) {
    // 移除自动生成的_id字段
    document.remove("_id");
    // 获取剩余字段的键值对
    String key = document.getFirstKey();
    Integer value = document.getInt64(key).intValue();
    return new Tuple2<>(key, value);
}

方案2:遍历字段排除_id并校验类型

通过遍历文档所有字段,跳过_id后提取有效整数字段:

@Override
public Tuple2<String, Integer> deserialize(BsonDocument document) {
    for (Map.Entry<String, BsonValue> entry : document.entrySet()) {
        String key = entry.getKey();
        // 跳过_id字段
        if (!"_id".equals(key)) {
            BsonValue value = entry.getValue();
            // 校验字段类型为INT64
            if (value.isInt64()) {
                return new Tuple2<>(key, value.asInt64().intValue());
            }
        }
    }
    throw new IllegalArgumentException("文档中未找到有效整数字段");
}

方案3:Sink阶段自定义_id(可选)

如果允许修改文档结构,可以在写入时将_id设为整数类型,避免后续混淆:

.setSerializationSchema(
        (input, context) 
            -> {
                Document doc = new Document(input.f0, input.f1);
                // 自定义_id为输入的整数值
                doc.put("_id", input.f1);
                return new InsertOneModel<>(BsonDocument.parse(doc.toJson()));
            })

内容的提问来源于stack exchange,提问作者Joseph Hwang

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 21:59:55