kafka消息中间件Java API调用

发布时间:2026/9/14 20:05:09

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(); } } } }消息消费的结果如下
延伸阅读

更多相关文章

2026/9/7 17:25:16

使用Tampermonkey油猴子给浏览器开个挂

油猴简介 一、油猴是什么 油猴(Tampermonkey)是免费的浏览器扩展和用户脚本管理器, 油猴子很特别, 它本身是一个无限手套(脚本管理器), 通过安装无限宝石(脚本), 能为我们提供超神的功能!它可以应用在多款浏览器上,比如谷歌浏览器,QQ浏览器&#xff0c…

2026/9/13 9:53:20

Linux系统构建高可用水产养殖监控集群实战

1. Linux系统龙虾部署:从零构建高可用集群的实战指南 在海鲜批发市场的数字化改造浪潮中,我遇到了一个有趣的挑战——如何用Linux系统构建稳定可靠的龙虾养殖环境监控平台。这个被我们戏称为"Linux系统龙虾部署"的项目,本质上是通过…

2026/9/14 20:00:23

技术项目命名指南:从无标题到好标题的实践

1. 项目概述作为一名从业多年的技术博主,我经常遇到这样的情况:手头有个不错的项目想法,却苦于找不到合适的标题来概括。这种情况在技术分享领域尤为常见——我们可能花了几周时间完成一个精彩的项目,却在最后一步"取名"…

2026/9/14 20:00:23

MATLAB图像处理在秸秆覆盖率估算中的应用

1. 秸秆覆盖率估算研究的背景与意义秸秆作为农业生产的重要副产品,其覆盖状况直接影响土壤质量、水分保持和作物生长。传统的人工测量方法存在效率低、主观性强等问题,而基于MATLAB的图像处理方法为解决这一难题提供了新的技术路径。在农业遥感领域&…

2026/9/14 20:00:23

COMSOL多物理场耦合在光伏集热器建模中的应用

1. 光伏集热器建模的独特挑战光伏集热器(PV-T)这个玩意儿确实挺有意思的,它把光伏发电和太阳能集热两个功能集成在一起,听起来很美好,但建模的时候简直就是个"混世魔王"。我去年给一家新能源企业做咨询时就遇…

2026/9/14 20:00:23

C++编译流程详解:从预处理到链接的完整指南

1. C编译方法概述:从源码到可执行文件的旅程刚接触C的新手往往会被这样的场景困扰:在终端输入g main.cpp后,一个可执行文件就神奇地出现了。但当你需要调试复杂项目时,这种"一步到位"的编译方式反而会成为效率杀手。实际…

2026/9/14 19:55:22

企业微信多账号接口实战:实例隔离与统一网关

「企业微信多账号接口」要解决的是:多个企微号同时运营多批外部群,数据不串、权限不混、掉线互不影响。 这篇讲接口层怎么做。 多账号模型 每个企微号一个 instance_id。所有登录、发送、回执、日志必须带它。账号绑定用途:推送号、接待号、…

2026/9/14 2:17:50

拯救者Y7000黑屏故障排查与维修实战指南

1. 项目概述:一台黑屏的拯救者Y7000,到底卡在哪一步? 联想拯救者Y7000系列笔记本,从2018年第一代搭载i5-8300H开始,到后来的i7-9750H、i7-10750H、i5-11400H,再到2023年款的R7-7840HS,它始终是学…

2026/9/14 0:03:22

KCF目标跟踪算法与OTB工程实现:毕业设计实战解析

简介:这是一份基于KCF核相关滤波算法、融合尺度池与抗遮挡处理的目标检测跟踪MATLAB完整源码,主要面向计算机相关专业准备毕业设计、课程设计或期末大作业的学生,也适合需要项目实战练习的初学者。源码在OTB数据集上完成验证,能够…

2026/9/14 0:03:22

语音情感识别实战:Keras实现LSTM、CNN、SVM与MLP多模型对比

简介:面向语音情感识别入门与进阶开发者,这份基于Keras的项目源码完整实现了LSTM、CNN、SVM、MLP四种模型,兼容Python3.8与Keras/TensorFlow2环境。压缩包内含49个文件,大小约70.31MB,主体包括Python脚本、yaml/json配…

2026/9/14 11:59:31

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/14 13:53:59

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/14 11:22:57

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码