kafka消息中间件Java API调用

发布时间:2026/9/24 8:06:52
kafka消息中间件Java API调用 参考: https://www.orchome.com/451Kafka集群的安装见上文本文介绍使用Java API通过kafka发送和接收消息。1. kafka客户端依赖dependency groupIdorg.apache.kafka/groupId artifactIdkafka_2.11/artifactId version1.0.1/version /dependency dependency groupIdorg.apache.kafka/groupId artifactIdkafka-clients/artifactId version1.0.1/version /dependency2 Kafka消息生产者APIpackage kafka; import java.util.Properties; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerRecord; public class ProducerDemo { private static final String MY_TOPIC my-topic; public static void main(String[] args) { Properties properties new Properties(); // Kafka 服务器地址 properties.put(bootstrap.servers, 127.0.0.1:9092,127.0.0.1:9093); // 消息应答机制 properties.put(acks, all); // 如果请求失败生产者会自动重试我们指定是0次如果启用重试则会有重复消息的可能性 properties.put(retries, 0); properties.put(batch.size, 16384); // 默认缓冲可立即发送即便缓冲空间还没有满但是如果你想减少请求的数量可以设置linger.ms大于0 properties.put(linger.ms, 1); // 控制生产者可用的缓存总量如果消息发送速度比其传输到服务器的快将会耗尽这个缓存空间 properties.put(buffer.memory, 33554432); // 消息序列化和反序列化方法 properties.put(key.serializer, org.apache.kafka.common.serialization.StringSerializer); properties.put(value.serializer, org.apache.kafka.common.serialization.StringSerializer); // 创建并发送消息 try (ProducerString, String producer new KafkaProducer(properties)) { for (int i 0; i 100; i) { String msg Message-index- i; producer.send(new ProducerRecord(MY_TOPIC, msg)); System.out.println(Sent: msg); } } } }消息发送结果:3. Kafka消息消费者APIpackage kafka; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import java.util.Collections; import java.util.Properties; public class ConsumerDemo { private static final String MY_TOPIC my-topic; public static void main(String[] args) { Properties properties new Properties(); // kafka 服务器地址 properties.put(bootstrap.servers, 127.0.0.1:9092); // 当前消费者所在的consumer group properties.put(group.id, group-1); // 消息消费后自动提交也可改为手动提交 properties.put(enable.auto.commit, true); // 自动提交间隔时间 properties.put(auto.commit.interval.ms, 1000); properties.put(auto.offset.reset, earliest); // 停止心跳的时间超过session.timeout.ms,那么就会认为是故障的它的分区将被分配到别的进程 properties.put(session.timeout.ms, 30000); // 消息序列化和反序列化方法 properties.put(key.deserializer, org.apache.kafka.common.serialization.StringDeserializer); properties.put(value.deserializer, org.apache.kafka.common.serialization.StringDeserializer); // 订阅 my-topic 主题的消息 KafkaConsumerString, String consumer new KafkaConsumer(properties); consumer.subscribe(Collections.singletonList(MY_TOPIC)); // 不停的获取消息并消费 while (true) { ConsumerRecordsString, String records consumer.poll(1000); System.out.println(records count: records.count()); for (ConsumerRecordString, String record : records) { System.out.printf(offset %d, key %s, value %s, record.offset(), record.key(), record.value()); System.out.println(); } } } }消息消费的结果如下

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询