关于Kafka 0.10 REST Proxy的三项技术问题咨询
1. Can two consumer instances exist in the same consumer group?
Absolutely—this is actually a standard scaling pattern for Kafka consumption, and the REST Proxy adheres to Kafka's core consumer group semantics. When multiple consumer instances belong to the same group, Kafka will evenly distribute the topic's partitions across all active instances in the group.
A quick heads-up: If the number of consumer instances in the group exceeds the number of partitions in the target topic, some instances will end up idle. That's because each partition can only be assigned to one consumer in a group at any given time.
2. What's the timeout period for a consumer instance after creation or last use?
In Kafka 0.10 REST Proxy, this timeout is controlled by the consumer.instance.timeout.ms configuration parameter. The default value is 300000 milliseconds (5 minutes).
If a consumer instance doesn’t send any requests (like polls, offset commits, or keep-alive calls) within this window, the REST Proxy will mark it as inactive and remove it from the consumer group. This frees up its assigned partitions to be reallocated to other active consumers in the group.
3. Are there Java or Python APIs for calling the Kafka REST API?
There’s no official, dedicated SDK for the Kafka REST API from Confluent or Apache for these languages, but you can easily use generic HTTP client libraries to interact with the REST endpoints directly:
Python Example
Using the popular requests library to create a consumer instance and poll for messages:
import requests # Base URL of your Kafka REST Proxy base_url = "http://your-rest-proxy-host:8082" # Create a consumer instance in a target group create_consumer = requests.post( f"{base_url}/consumers/my-group", json={ "name": "my-consumer-instance", "format": "json", "auto.offset.reset": "latest" } ) consumer_instance = create_consumer.json() # Poll for messages from a topic poll_response = requests.get( f"{base_url}/consumers/my-group/instances/my-consumer-instance/topics/my-topic" ) messages = poll_response.json() print(messages)
Java Example
Using Apache HttpClient to perform a similar poll operation:
import org.apache.http.client.fluent.Request; import org.apache.http.entity.ContentType; import java.io.IOException; public class KafkaRestClient { public static void main(String[] args) throws IOException { String baseUrl = "http://your-rest-proxy-host:8082"; // Create consumer instance String createConsumerResponse = Request.Post(baseUrl + "/consumers/my-group") .bodyString("{\"name\":\"my-consumer-instance\",\"format\":\"json\",\"auto.offset.reset\":\"latest\"}", ContentType.APPLICATION_JSON) .execute() .returnContent() .asString(); System.out.println(createConsumerResponse); // Poll for messages String pollResponse = Request.Get(baseUrl + "/consumers/my-group/instances/my-consumer-instance/topics/my-topic") .execute() .returnContent() .asString(); System.out.println(pollResponse); } }
These are basic examples—you can extend them to handle offset commits, consumer instance deletion, and other operations supported by the Kafka REST API.
内容的提问来源于stack exchange,提问作者Jack

