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

HTTP Sink Connector未按配置批量发送消息,求缺失配置建议

HTTP Sink Connector批量发送问题排查

你的配置里已经开启了batching.enabled=true和batch.max.size=50,但仍逐条发送,主要缺少以下关键配置项:

  • batch.json.as.array=true:当request.body.format=json时,默认不会将批量消息包装为JSON数组,而是会逐条发送。添加这个配置后,连接器会把最多50条消息组装成一个JSON数组作为请求体发送。
  • batch.delay.ms(可选但建议):设置超时时间(例如1000,单位毫秒),当消息数量未达到batch.max.size但超过该时间时,也会发送当前攒下的批量消息,避免因消息量不足一直等待。

另外需要确认两点:

  • 目标HTTP端点需支持接收JSON数组格式的请求体,否则批量发送会返回错误,连接器可能退化为逐条重试。
  • 检查Kafka Connect的consumer.max.poll.records(全局配置)或连接器级的consumer.override.max.poll.records,该值需大于等于batch.max.size,否则连接器一次拉取的消息数不足,无法凑够批量。

修正后的配置示例(添加关键项):

{
"name": "HTTPSinkConnector_1",
"config": {
    "topics": "topic_1",
    "tasks.max": "1",
    "connector.class": "io.confluent.connect.http.HttpSinkConnector",
    "http.api.url": "http://localhost/messageHandler",
    "request.method": "POST",
    "key.converter":"org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "io.confluent.connect.avro.AvroConverter",
    "value.converter.schema.registry.url": "http://schema-registry:8081",
    "confluent.topic.bootstrap.servers": "kafka:19092",
    "confluent.topic.replication.factor": "1",
    "batching.enabled": true,
    "batch.max.size": 50,
    "batch.json.as.array": true,
    "batch.delay.ms": 1000,
    "consumer.override.max.poll.records": 50,
    "reporter.bootstrap.servers": "kafka:19092",
    "reporter.result.topic.name": "success-responses",
    "reporter.result.topic.replication.factor": "1",
    "reporter.error.topic.name": "error-responses",
    "reporter.error.topic.replication.factor": "1",
    "request.body.format": "json"
 }
}

内容的提问来源于stack exchange,提问作者Dixit Singla

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 10:45:25