如何在Flink Java MongoSink中设置动态集合名?
解决Flink写入MongoDB时动态指定集合名的问题
要实现根据ObjectFromKafka中的collectionName字段动态设置MongoDB集合名,只需要替换掉静态的集合名配置,改用函数式的集合名指定方式即可,不用开发复杂的自定义Sink。
修改后的完整代码
KafkaSource<ObjectFromKafka> kafkaSource = KafkaSource.<ObjectFromKafka>builder() // .setTopics(properties.getProperty(TOPIC_READ)) // .setProperties(properties) // .setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST)) // .setValueOnlyDeserializer(new JsonDeserializationSchema<>(ObjectFromKafka.class)) // .build(); DataStream<ObjectFromKafka> stream = env.fromSource(kafkaSource, WatermarkStrategy.forMonotonousTimestamps(), "kafka_source"); MongoSerializationSchema<ObjectFromKafka> serializationSchema = (input, context) -> new InsertOneModel<>(BsonDocument.parse(convertToJsonString(input))); MongoSink<ObjectFromKafka> sink1 = MongoSink.<ObjectFromKafka>builder() // .setUri("mongodb://username:password@mongo:27017/") // .setDatabase("admin") // // 关键修改:用动态函数替换静态集合名 .setCollection(input -> input.getCollectionName()) .setBatchSize(1000) // .setBatchIntervalMs(1000) // .setMaxRetries(3) // .setDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE) // .setSerializationSchema(serializationSchema) // .build(); stream.sinkTo(sink1).name("mongoDb Sink");
关键说明
- 动态集合名的核心:
setCollection方法支持传入一个Function<T, String>类型的参数,这个函数会针对每条输入数据执行,从当前的ObjectFromKafka实例中提取collectionName字段值,作为该条数据要写入的MongoDB集合名称。 - 异常情况处理:如果担心
collectionName字段可能为空或无效,可以在函数中添加兜底逻辑,比如:.setCollection(input -> { String colName = input.getCollectionName(); return colName != null && !colName.isEmpty() ? colName : "default_collection"; }) - 前提条件:确保你的
ObjectFromKafka类已经定义了collectionName字段对应的getter方法getCollectionName()。
内容的提问来源于stack exchange,提问作者Domix
相关产品推荐
相关产品推荐

