如何在Kafka 0.10 REST API中指定消费者偏移量及文档咨询
Kafka 0.10 REST API:指定消费偏移量及完整参数指南
我之前折腾过Kafka 0.10的REST API,刚好碰到过和你一样的困惑——想指定偏移量消费但找不到参数,给你详细说下解决方案和完整的参数说明:
一、如何指定消费偏移量?
Kafka 0.10的REST代理(Confluent维护的官方代理)并没有直接在拉取消息的接口里提供偏移量参数,得通过「先设置偏移量,再拉取消息」的两步操作实现:
- 创建消费者实例
先发送POST请求到/consumers/{你的消费组名},创建一个消费者实例。请求体可以带上基础配置,比如:
{ "auto.commit.enable": "false", "format": "json" }
(这里建议关闭自动提交,避免偏移量被自动覆盖)
- 手动设置目标偏移量
接着发送POST请求到/consumers/{你的消费组名}/instances/{消费者实例名}/offsets,请求体里指定要消费的topic、分区和具体偏移量:
{ "offsets": [ { "topic": "your-topic", "partition": 0, "offset": 12345 } ] }
- 拉取消息
最后发送GET请求到/consumers/{你的消费组名}/instances/{消费者实例名}/records,就能从你指定的偏移量开始消费了。
二、完整的REST代理核心参数说明
下面整理了Kafka 0.10版本REST代理的关键接口参数,覆盖消费和生产常用场景:
消费者相关接口
1. 创建消费者(POST /consumers/{group_name})
请求体参数:
name: 可选,消费者实例的自定义名称,不填会自动生成唯一IDformat: 可选,消息序列化格式,支持json、avro、binary,默认jsonauto.offset.reset: 可选,无初始偏移量或偏移量无效时的 fallback 策略,可选earliest(从头读)、latest(读最新),默认latestauto.commit.enable: 可选,是否自动提交消费偏移量,默认trueauto.commit.interval.ms: 可选,自动提交的时间间隔,默认5000msfetch.min.bytes: 可选,拉取消息的最小字节数,默认1字节(凑够字节数才返回)fetch.wait.max.ms: 可选,拉取的最长等待时间,默认500ms(到点不管够不够字节都返回)
2. 设置偏移量(POST /consumers/{group_name}/instances/{consumer_instance}/offsets)
请求体的offsets数组每个元素参数:
topic: 必选,目标topic名称partition: 必选,分区编号offset: 必选,要设置的具体偏移量值metadata: 可选,偏移量的备注信息(比如设置原因)
3. 拉取消息(GET /consumers/{group_name}/instances/{consumer_instance}/records)
URL查询参数:
timeout: 可选,拉取超时时间,默认500msmax_bytes: 可选,单次拉取的最大字节数,默认1MBtopics: 可选,要拉取的topic列表(如果创建消费者时没订阅topic,这里必须传)
生产者相关接口(顺带整理了常用的)
发送消息(POST /topics/{topic_name})
请求体参数:
records: 必选,消息数组,每个元素包含:key: 可选,消息键(用于分区路由),格式和format参数对应value: 必选,消息内容partition: 可选,指定发送到的分区编号,不填则按key哈希分配
key_schema_id: 可选,使用avro格式时的key schema IDvalue_schema_id: 可选,使用avro格式时的value schema ID
三、Java HTTP调用的注意事项
用Java的HTTP客户端(比如OkHttp、Apache HttpClient)调用时,一定要设置正确的请求头:
- 当使用json格式时,Content-Type为
application/vnd.kafka.json.v2+json - avro格式对应
application/vnd.kafka.avro.v2+json - binary格式对应
application/vnd.kafka.binary.v2+json
内容的提问来源于stack exchange,提问作者Jack
相关产品推荐
相关产品推荐

