如何配置Confluent HttpSink Connector设置Content-Type为application/json
问题描述
已搭建本地Kafka集群与微服务端点,希望HttpSink Connector接收主题中类似{"foo":"bar"}的消息,并以JSON字符串形式调用API端点。未使用Schema Registry,当前连接器配置如下:
{ "name": "httpsink3", "config": { "name": "httpsink3", "connector.class": "io.confluent.connect.http.HttpSinkConnector", "tasks.max": "1", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter": "org.apache.kafka.connect.storage.StringConverter", "topics": "test3", "http.api.url": "http://localhost:5000/out", "reporter.result.topic.name": "success-responses", "reporter.result.topic.replication.factor": "1", "reporter.error.topic.name": "error-responses", "reporter.error.topic.replication.factor": "1", "reporter.bootstrap.servers": "localhost:9092", "confluent.topic.bootstrap.servers": "localhost:9092", "confluent.topic.replication.factor": "1" } }
实际调用结果显示请求头Content-Type为text/plain; charset=UTF-8,期望设置为application/json。尝试使用"value.converter": "org.apache.kafka.connect.json.JsonConverter"后,请求头仍为text/plain,且消息内容变为{"foo"="bar"},格式异常。请问该如何配置才能满足需求?
解决方案
核心配置调整
保持value.converter为StringConverter(因为主题中的消息已经是JSON格式的字符串,无需额外转换),同时添加http.headers.content.type参数强制设置请求头为application/json,即可解决问题。
完整配置示例
{ "name": "httpsink3", "config": { "name": "httpsink3", "connector.class": "io.confluent.connect.http.HttpSinkConnector", "tasks.max": "1", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter": "org.apache.kafka.connect.storage.StringConverter", "topics": "test3", "http.api.url": "http://localhost:5000/out", // 新增:强制设置请求Content-Type为application/json "http.headers.content.type": "application/json", "reporter.result.topic.name": "success-responses", "reporter.result.topic.replication.factor": "1", "reporter.error.topic.name": "error-responses", "reporter.error.topic.replication.factor": "1", "reporter.bootstrap.servers": "localhost:9092", "confluent.topic.bootstrap.servers": "localhost:9092", "confluent.topic.replication.factor": "1" } }
配置说明
value.converter: org.apache.kafka.connect.storage.StringConverter:主题中的消息本身就是JSON格式的字符串,使用StringConverter可以直接将消息内容原封不动地传递给API,避免格式转换错误。http.headers.content.type: "application/json":HttpSink Connector默认会根据转换器类型设置Content-Type,用StringConverter时默认是text/plain,通过该参数可以强制覆盖为application/json,符合API的接收要求。- 之前使用
JsonConverter出现格式异常的原因:未配置Schema Registry时,JsonConverter会将结构化数据转换为一种带等号的非标准格式(类似{"foo"="bar"}),这与API期望的标准JSON格式不符,因此不适合当前场景。
内容的提问来源于stack exchange,提问作者FrankZhu
相关产品推荐
相关产品推荐

