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

优化Flink流作业中自定义MongoDB Sink的写入性能

Flink自定义MongoDB Sink写入性能优化方案

问题背景

我运行着一个Kafka->Map->MongoDB的Flink流作业,由于没有适配Flink的MongoDB Sink连接器,因此自定义了一个用于向MongoDB写入数据的Sink。然而写入操作的吞吐量较低,尽管尝试使用MongoDB Reactive Streams驱动实现异步操作,但结果仍未达到业务需求(仅约2000次写入/秒,而业务需要至少5000次写入/秒)。

当前自定义Sink代码

import com.mongodb.WriteConcern;
import com.mongodb.client.result.InsertOneResult;
import com.mongodb.reactivestreams.client.MongoClient;
import com.mongodb.reactivestreams.client.MongoClients;
import com.mongodb.reactivestreams.client.MongoDatabase;
import org.apache.flink.configuration.ConfigOptions;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
import org.bson.Document;
import org.reactivestreams.Publisher;
import org.reactivestreams.Subscriber;
import org.reactivestreams.Subscription;

import java.util.logging.Logger;

public class MongoDBSinkConfig extends RichSinkFunction<Document> {

    Configuration configuration;

    static MongoClient mongoClient;

    static MongoDatabase mongoDatabase;

    public MongoDBSinkConfig(Configuration configuration) {
        this.configuration = configuration;
    }

    @Override
    public void open(Configuration parameters) {
        String uri = configuration.getString(ConfigOptions.key(ConfigurationParamUtils.MONGO_URI).stringType().noDefaultValue());
        String database = configuration.getString(ConfigOptions.key(ConfigurationParamUtils.MONGO_DATABASE).stringType().noDefaultValue());
        String config = configuration.getString(ConfigOptions.key(ConfigurationParamUtils.CONFIG_COLLECTION).stringType().noDefaultValue());
        MongoDBMetricConfig.loadMongoDBMetricConfig(uri, database, config);
        mongoClient = MongoClients.create(uri);
        mongoDatabase = mongoClient.getDatabase(database).withWriteConcern(WriteConcern.UNACKNOWLEDGED);
    }

    @Override
    public void invoke(Document value, Context context) {
        Publisher<InsertOneResult> insertResult = mongoDatabase.getCollection("collection").withWriteConcern(WriteConcern.UNACKNOWLEDGED).insertOne(value);
        
        insertResult.subscribe(new Subscriber<InsertOneResult>() {
            @Override
            public void onSubscribe(Subscription s) {

            }

            @Override
            public void onNext(InsertOneResult insertOneResult) {
                Logger.getGlobal().info(insertOneResult.toString());
            }

            @Override
            public void onError(Throwable t) {
                Logger.getGlobal().info(t.getMessage());
            }

            @Override
            public void onComplete() {
            }
        });
    }

    @Override
    public void close() {
        SinkUtils.closeConnection();
    }
}

性能优化方法

  • 改用批量写入替代单条插入:单条insertOne会产生大量网络往返请求,是吞吐量瓶颈的核心原因。可以在Sink中维护本地缓存队列,当队列达到指定大小(如1000条)或超时时间时,调用insertMany批量写入,大幅减少网络开销。
  • 修复Reactive Streams订阅逻辑:当前Subscriber的onSubscribe方法为空,未调用s.request(Long.MAX_VALUE)请求数据,会导致Publisher无法主动推送结果,阻塞异步流程。必须添加该调用,确保异步操作正常流转。
  • 复用MongoCollection实例:当前每次invoke都调用getCollection("collection")会重复创建资源,建议在open方法中初始化并缓存MongoCollection<Document>实例,后续直接复用。
  • 优化MongoDB客户端配置:通过MongoClientSettings调整连接池参数,比如增大maxPoolSize(默认100,可根据并发需求调整至200-500)、设置合理的socketTimeout和connectTimeout,避免连接等待损耗性能。示例:
    MongoClientSettings settings = MongoClientSettings.builder()
        .applyConnectionString(new ConnectionString(uri))
        .applyToConnectionPoolSettings(builder -> builder.maxSize(300))
        .applyToSocketSettings(builder -> builder.readTimeout(Duration.ofSeconds(10)))
        .build();
    mongoClient = MongoClients.create(settings);
    
  • 调整Flink作业并行度:Sink的并行度应与上游算子匹配,同时不超过MongoDB集群承受的并发连接上限。通过env.setParallelism()或Sink单独设置setParallelism()提升并发写入能力。
  • 移除同步日志输出:代码中onNext和onError使用的同步日志会阻塞异步写入流程,建议改用异步日志框架(如Logback异步appender),或直接移除非必要日志。
  • 使用Flink AsyncFunction实现异步写入:相比自定义RichSink的异步实现,Flink的AsyncFunction可更好控制异步请求并发度(通过setCapacity),避免请求过载,同时更贴合Flink流处理模型。
  • 优化MongoDB集群配置:检查索引是否冗余,过多索引会降低写入速度;分片集群需确保分片键合理,避免热点分片;同时确认集群硬件资源(磁盘IO、CPU、内存)无瓶颈,比如使用SSD提升写入性能。
  • 谨慎调整WriteConcern:当前已用WriteConcern.UNACKNOWLEDGED(最低确认级别),如果业务允许可保持;若需一定可靠性,可考虑WriteConcern.W1(仅确认主节点写入),平衡可靠性与性能。

内容的提问来源于stack exchange,提问作者Lannest Caster

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 03:20:33