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
相关产品推荐
相关产品推荐

