Kafka HTTP Sink Connector报415错误,请求排查配置问题
Kafka Connect HTTP Sink 415错误排查与.NET API接收方案
一、HTTP Sink Connector配置问题排查
415错误核心是请求媒体类型与API支持的类型不匹配,优先检查以下配置项:
1. 配置正确的Content-Type与数据格式
默认情况下,HTTP Sink Connector会发送原始Avro二进制(application/octet-stream),但你的.NET API大概率期望JSON格式。需添加以下配置:
http.content.type=application/json:指定请求的媒体类型为JSONformat=json:让连接器自动将Avro数据转换为JSON结构发送
2. 验证Avro转换器配置
因为Topic存储的是Avro格式数据,必须配置Avro转换器对接Schema Registry,确保能正确解析Topic中的数据:
{ "name": "your-http-sink", "config": { "connector.class": "io.confluent.connect.http.HttpSinkConnector", "tasks.max": "1", "topics": "target-topic", "http.api.url": "http://your-dotnet-api/your-endpoint", "http.content.type": "application/json", "format": "json", "value.converter": "io.confluent.connect.avro.AvroConverter", "value.converter.schema.registry.url": "http://your-schema-registry:8081", "key.converter": "org.apache.kafka.connect.storage.StringConverter" } }
3. 排除其他配置问题
- 确认Schema Registry地址正确,连接器能拉取到对应Topic的Avro Schema
- 检查连接器是否有权限访问目标API(curl测试正常的话,大概率不是此问题)
二、.NET API正确接收消息的实现
根据连接器发送的格式,选择对应的接收方式:
方式1:接收JSON格式(推荐)
如果连接器已配置为发送JSON,直接用avrogen生成的DTO作为接收参数,只需确保JSON序列化配置匹配字段命名:
using Microsoft.AspNetCore.Mvc; [ApiController] [Route("api/[controller]")] public class MessageController : ControllerBase { [HttpPost] [Consumes("application/json")] public IActionResult Receive([FromBody] YourAvroGeneratedDto message) { // 处理业务逻辑 return Ok(); } }
字段命名匹配配置
如果Avro Schema使用蛇形命名(如user_name),而avrogen生成的C#属性为驼峰(如UserName),需在Startup/Program.cs中配置JSON序列化策略:
// System.Text.Json示例 builder.Services.AddControllers() .AddJsonOptions(options => { options.JsonSerializerOptions.PropertyNamingPolicy = JsonNamingPolicy.SnakeCaseLower; options.JsonSerializerOptions.PropertyNameCaseInsensitive = true; });
方式2:接收原始Avro二进制(仅特殊场景使用)
若必须发送Avro二进制,需修改API端点解析字节流:
using Avro.IO; using Microsoft.AspNetCore.Mvc; [ApiController] [Route("api/[controller]")] public class AvroMessageController : ControllerBase { [HttpPost] [Consumes("application/octet-stream")] public async Task<IActionResult> ReceiveAvro() { using var stream = new MemoryStream(); await Request.Body.CopyToAsync(stream); stream.Position = 0; // 使用avrogen生成的Schema创建读取器 var reader = DatumReader<YourAvroGeneratedDto>.CreateDefaultReader(YourAvroGeneratedDto.Schema); var decoder = new BinaryDecoder(stream); var message = reader.Read(null, decoder); // 处理业务逻辑 return Ok(); } }
此时连接器需配置http.content.type=application/octet-stream,且移除format=json配置。
内容的提问来源于stack exchange,提问作者Dixit Singla
相关产品推荐
相关产品推荐

