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

