【发布时间】:2017-07-03 00:05:43
【问题描述】:
我有一个使用kafka client 2.11: 0.10.2.1 的小型Java spark 服务。
以下是当我阅读从最新 Kafka 版本发布的主题时运行良好的代码:
Properties props = new Properties();
props.put(org.apache.kafka.clients.producer.ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, producerConfig.getBrokerConnectionString());
props.put(org.apache.kafka.clients.producer.ProducerConfig.ACKS_CONFIG, "all");
props.put(org.apache.kafka.clients.producer.ProducerConfig.RETRIES_CONFIG, producerConfig.getRetry());
props.put(org.apache.kafka.clients.producer.ProducerConfig.BATCH_SIZE_CONFIG, producerConfig.getBatchSize());
props.put(org.apache.kafka.clients.producer.ProducerConfig.LINGER_MS_CONFIG, producerConfig.getLingerTimeInMs());
props.put(org.apache.kafka.clients.producer.ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, producerConfig.getRequestTimeout());
props.put(org.apache.kafka.clients.producer.ProducerConfig.MAX_BLOCK_MS_CONFIG, producerConfig.getMaxBlockMS());
props.put(org.apache.kafka.clients.producer.ProducerConfig.CONNECTIONS_MAX_IDLE_MS_CONFIG, producerConfig.getMaxIdleTime());
props.put(org.apache.kafka.clients.producer.ProducerConfig.BUFFER_MEMORY_CONFIG, maxBytesInBuffer / producerConfig.getProducersCount());
props.put(org.apache.kafka.clients.producer.ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
props.put(org.apache.kafka.clients.producer.ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer");
producers = new Producer[1];
producers[0] = new KafkaProducer<>(props);
producers[0].partitionsFor("mYTopic").size();
已有一个Kafka主题,其中kafka版本为0.8.2.x
.我也想为此使用相同的代码。但是这段代码在最后一行(partitionsFor)中给出了超时,主题是 Kafka 发布的 0.8.2.x 版本。在这方面的任何帮助将不胜感激。
简而言之:Kafka 主题(由 0.8.2.x 发布)无法被 0.10.2.1 客户端读取
【问题讨论】:
-
文件 kafka.apache.org/documentation/#upgrade 说 0.10.2 客户无法与 0.8.2 经纪人交谈。
标签: java apache-kafka kafka-consumer-api