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

如何配置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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 13:30:47