资讯中心

Kafka 消息中间件

📅 2026/8/13 14:59:00
Kafka 消息中间件
一消息中间件消息中间件是在消息的传输过程中保存消息的容器。消息中间件在将消息从消息生产者到消费者时充当中间人的作用。队列的主要目的是提供路由并保证消息的传送如果发送消息时接收者不可用消息对列会保留消息直到可以成功地传递它为止当然消息队列保存消息也是有期限的。二消息中间件特点1解耦允许你独立的扩展或修改两边的处理过程只要确保它们遵守同样的接口约束。2冗余消息队列把数据进行持久化直到它们已经被完全处理通过这一方式规避了数据丢失风险。许多消息队列所采用的插入-获取-删除范式中在把一个消息从队列中删除之前需要你的处理系统明确的指出该消息已经被处理完毕从而确保你的数据被安全的保存直到你使用完毕。3扩展性因为消息队列解耦了你的处理过程所以增大消息入队和处理的频率是很容易的只要另外增加处理过程即可。4消峰在访问量剧增的情况下应用仍然需要继续发挥作用但是这样的突发流量并不常见。如果为以能处理这类峰值访问为标准来投入资源随时待命无疑是巨大的浪费。使用消息队列能够使关键组件顶住突发的访问压力而不会因为突发的超负荷的请求而完全崩溃。5可恢复性系统的一部分组件失效时不会影响到整个系统。消息队列降低了进程间的耦合度所以即使一个处理消息的进程挂掉加入队列中的消息仍然可以在系统恢复后被处理。6顺序保证在大多使用场景下数据处理的顺序都很重要。大部分消息队列本来就是排序的并且能保证数据会按照特定的顺序来处理。Kafka保证一个Partition内的消息的有序性7缓冲有助于控制和优化数据流经过系统的速度解决生产消息和消费消息的处理速度不一致的情况。8异步通信很多时候用户不想也不需要立即处理消息。消息队列提供了异步处理机制允许用户把一个消息放入队列但并不立即处理它。想向队列中放入多少消息就放多少然后在需要的时候再去处理它们。消息模型在JMS标准中有两种消息模型点对点Point to Point,发布/订阅(Pub/Sub)。P2P模式P2P模式包含三个角色消息队列Queue发送者(Sender)接收者(Receiver)。消息生产者生产消息发送到queue中然后消息消费者从queue中取出并且消费消息。消息被消费以后queue中不再有存储所以消息消费者不可能消费到已经被消费的消息。Queue支持存在多个消费者但是对一个消息而言只会有一个消费者可以消费。Pub/sub模式包含三个角色主题Topic发布者Publisher订阅者Subscriber。消息生产者发布将消息发布到topic中同时有多个消息消费者订阅消费该消息。和点对点方式不同发布到topic的消息会被所有订阅者消费。Kafka是一种分布式消息系统由LinkedIn使用Scala编写用作LinkedIn的活动流(Activity Stream)和运营数据处理管道(Pipeline)的基础具有高水平扩展和高吞吐量。目前越来越多的开源分布式处理系统如Apache flume、Apache Storm、Spark、Elasticsearch都支持与Kafka集成。 产品基于高可用分布式集群技术提供消息订阅和发布、消息轨迹查询、定时延时消息、资源统计、监控报警等一系列消息云服务是企业级互联网架构的核心产品。首先打个比方kafka好比就是电视台而电视台下面有很多节目生产者就是制作节目的团队而消费者就是我们观看这个节目的人一开始在zookeeper创建一个节目假设就叫cctv1有了这个节目名后我们就得请一个团队来填充这个节目比如放电视剧之类的数据而我们消费者要观看这个节目的话就得需要zookeeper来授权给我们。电视台则只是存数据的相当于一个中间人和现在中介差不多个意思kafka基本概念producer生产者发布消息到 kafka 集群的终端或服务。consumer消费者从 kafka 集群中消费消息的终端或服务。topic: 消息以topic为类别记录,每一类的消息称之为一个主题(Topic)可以理解为一个队列。broker以集群的方式运行,可以由一个或多个服务组成每个服务叫做一个broker;消费者可以订阅一个或多个主题(topic), 并从Broker拉数据,从而消费这些已发布的消息。PartitionTopic物理上的分组一个topic可以分为多个partition每个partition是一个有序的队列。partition中的每条消息都会被分配一个有序的idoffset。Segmentpartition物理上由多个segment组成每个segment file对应两个文件分别是以.log结尾的数据文件和以.index结尾的索引文件。Message消息是通信的基本单位每个producer可以向一个topic主题发布一些消息。Consumer Group 每个Consumer属于一个特定的Consumer Group可为每个Consumer指定group name若不指定group name则属于默认的group。Offsetkafka的存储文件都是按照offset.kafka来命名用offset做名字的好处是方便查找。例如你想找位于2049的位置只要找到2048.kafka的文件即可。当然the first offset就是00000000000.kafka。消息发送的流程Producer根据指定的partition方法round-robin、hash等将消息发布到指定topic的partition里面kafka集群接收到Producer发过来的消息后将其持久化到硬盘并保留消息指定时长可配置而不关注消息是否被消费。Consumer从kafka集群pull数据并控制获取消息的offsetBroker保存消息存储方式物理上把topic分成一个或多个patition对应 server.properties 中的num.partitions3配置每个patition物理上对应一个文件夹partiton命名规则为topic名称有序序号第一个partiton序号从0开始序号最大值为partitions数量减1。[atguiguhadoop102 logs]$ ll drwxrwxr-x. 2 atguigu atguigu 4096 8月 6 14:37 first-0 drwxrwxr-x. 2 atguigu atguigu 4096 8月 6 14:35 first-1 drwxrwxr-x. 2 atguigu atguigu 4096 8月 6 14:37 first-2每个partion(目录)相当于一个巨型文件被平均分配到多个大小相等segment(段)数据文件中如下[atguiguhadoop102 logs]$ cd first-0 [atguiguhadoop102 first-0]$ ll -rw-rw-r--. 1 atguigu atguigu 10485760 8月 6 14:33 00000000000000000000.index -rw-rw-r--. 1 atguigu atguigu 219 8月 6 15:07 00000000000000000000.log -rw-rw-r--. 1 atguigu atguigu 10485756 8月 6 14:33 00000000000000000000.timeindex -rw-rw-r--. 1 atguigu atguigu 8 8月 6 14:37 leader-epoch-checkpoint存储策略无论消息是否被消费kafka都会保留所有消息。有两种策略可以删除旧数据基于时间log.retention.hours168基于文件大小log.retention.bytes1073741824需要注意的是因为Kafka读取特定消息的时间复杂度为O(1)即与文件大小无关所以这里删除过期文件与提高 Kafka 性能无关。Kafka消息发送方式Kafka消息发送分同步(sync)、异步(async)两种方式默认是使用同步方式可通过producer.type属性进行配置Kafka保证消息被安全生产有三个选项分别是0,1,-1通过request.required.acks属性进行配置0代表不进行消息接收是否成功的确认(默认值)1代表当Leader副本接收成功后返回接收成功确认信息-1代表当Leader和Follower副本都接收成功后返回接收成功确认信息六种发送场景消息丢失的场景网络异常acks设置为0时不和Kafka集群进行消息接受确认当网络发生异常等情况时存在消息丢失的可能客户端异常异步发送时消息并没有直接发送至Kafka集群而是在Client端按一定规则缓存并批量发送。在这期间如果客户端发生死机等情况都会导致消息的丢失缓冲区满了异步发送时Client端缓存的消息超出了缓冲池的大小也存在消息丢失的可能Leader副本异常acks设置为1时Leader副本接收成功Kafka集群就返回成功确认信息而Follower副本可能还在同步。这时Leader副本突然出现异常新Leader副本(原Follower副本)未能和其保持一致就会出现消息丢失的情况以上就是消息丢失的几种情况在日常应用中我们需要结合自身的应用场景来选择不同的配置。想要更高的吞吐量就设置异步、ack0想要不丢失消息数据就选同步、ack-1策略消息接收的三种模式kafka的消费模式总共有3种最多一次最少一次正好一次。为什么会有这3种模式是因为客户端处理消息提交反馈commit这两个动作不是原子性。最多一次客户端收到消息后在处理消息前自动提交这样kafka就认为consumer已经消费过了偏移量增加。最少一次客户端收到消息处理消息再提交反馈。这样就可能出现消息处理完了在提交反馈前网络中断或者程序挂了那么kafka认为这个消息还没有被consumer消费产生重复消息推送。正好一次保证消息处理和提交反馈在同一个事务中即有原子性。消费消息有没有顺序Kafka中分区只能被同一个消费者组中的一个消费者消费但是可以被多个不同消费者组的消费者同时消费。如果消费者组中只有一个消费者那么这个消费者可以消费所有的分区。既然kafka允许多个消费者对多个分区同时消费且生产者生产的消息也落于不同的分区中那么在这种情况下消费这些消息的顺序肯定是不可控的。我们可以控制生产消息的过程中让消息落入同一个分区通过设定消息的key让kafka生产者根据key进行hash选择要写入的分区来保证消息写入的顺序以及消费的顺序。kafka的消费模式Kafka的消费模式主要有两种一种是一对一的消费也即点对点的通信即一个发送一个接收。第二种为一对多(发布/订阅模式)的消费即一个消息发送到消息队列消费者根据消息队列的订阅拉取信息消费。发布/订阅模式即利用Topic存储消息消息生产者将消息发布到Topic中同时有多个消费者订阅此topic消费者可以从中消费消息注意发布到Topic中的消息会被多个消费者消费消费者消费数据之后数据不会被清除而是按照时间策略来删除Kafka会默认保留一段时间然后再删除。kafka的备份策略Kafka的备份的单元是partition也就是每个partition都会有leader partiton和follow partiton。其中leader partition是用来进行和producer进行写交互follow从leader副本进行拉数据进行同步从而保证数据的冗余防止数据丢失的目的。当 partition 对应的 leader 宕机时需要从 follower 中选举出新 leader。在选举新leader时一个基本的原则是新的 leader 必须拥有旧 leader commit 过的所有消息。Kafka的应用场景1、消息队列比起大多数的消息系统来说Kafka有更好的吞吐量内置的分区冗余及容错性这让Kafka成为了一个很好的大规模消息处理应用的解决方案。消息系统一般吞吐量相对较低但是需要更小的端到端延时并尝尝依赖于Kafka提供的强大的持久性保障。在这个领域Kafka足以媲美传统消息系统如ActiveMR或RabbitMQ。2、行为跟踪Kafka的另一个应用场景是跟踪用户浏览页面、搜索及其他行为以发布-订阅的模式实时记录到对应的topic里。那么这些结果被订阅者拿到后就可以做进一步的实时处理或实时监控或放到hadoop/离线数据仓库里处理。3、元信息监控作为操作记录的监控模块来使用即汇集记录一些操作信息可以理解为运维性质的数据监控吧。4、日志收集日志收集方面其实开源产品有很多包括Scribe、Apache Flume。很多人使用Kafka代替日志聚合log aggregation。日志聚合一般来说是从服务器上收集日志文件然后放到一个集中的位置文件服务器或HDFS进行处理。然而Kafka忽略掉文件的细节将其更清晰地抽象成一个个日志或事件的消息流。这就让Kafka处理过程延迟更低更容易支持多数据源和分布式数据处理。比起以日志为中心的系统比如Scribe或者Flume来说Kafka提供同样高效的性能和因为复制导致的更高的耐用性保证以及更低的端到端延迟。5、流处理这个场景可能比较多也很好理解。保存收集流数据以提供之后对接的Storm或其他流式计算框架进行处理。很多用户会将那些从原始topic来的数据进行阶段性处理汇总扩充或者以其他的方式转换到新的topic下再继续后面的处理。例如一个文章推荐的处理流程可能是先从RSS数据源中抓取文章的内容然后将其丢入一个叫做“文章”的topic中后续操作可能是需要对这个内容进行清理比如回复正常数据或者删除重复数据最后再将内容匹配的结果返还给用户。这就在一个独立的topic之外产生了一系列的实时数据处理的流程。Strom和Samza是非常著名的实现这种类型数据转换的框架。6、事件源事件源是一种应用程序设计的方式该方式的状态转移被记录为按时间顺序排序的记录序列。Kafka可以存储大量的日志数据这使得它成为一个对这种方式的应用来说绝佳的后台。比如动态汇总News feed。7、持久性日志commit logKafka可以为一种外部的持久性日志的分布式系统提供服务。这种日志可以在节点间备份数据并为故障节点数据回复提供一种重新同步的机制。Kafka中日志压缩功能为这种用法提供了条件。在这种用法中Kafka类似于Apache BookKeeper项目。Docker Desktop安装Kafka新建文件夹如d:\kafka-docker在里面建docker-compose.ymlversion: 3.8 services: kafka: image: docker.1ms.run/apache/kafka:3.8.1 container_name: kafka restart: always ports: - 9092:9092 - 9093:9093 environment: KAFKA_NODE_ID: 1 KAFKA_PROCESS_ROLES: broker,controller KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092 KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT KAFKA_CONTROLLER_QUORUM_VOTERS: 1kafka:9093 KAFKA_LOG_DIRS: /var/lib/kafka/data CLUSTER_ID: 06M1KKGfRuyK7XWurAYXkA KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 volumes: - ./data:/var/lib/kafka/data核心要点Kafka 3.x 起KRaft 模式已生产就绪完全不需要 ZooKeeper9092是业务端口9093是 KRaft 内部控制器通信端口整个服务只需1 个容器比 ZK 模式简洁太多localhost:9092适合本地开发外部访问需改为宿主机 IP启动kafka,docker会自动安装kakkadocker-compose up -d看到started就启动成功。[2026-05-28 02:53:20,084] INFO [KafkaRaftServer nodeId1] Kafka Server started (kafka.server.KafkaRaftServer)powershell命令进入容器docker exec -it kafka bash cd /opt/kafka创建 topicbin/kafka-topics.sh --create --topic test-topic --bootstrap-server localhost:9092消费者组订阅新topic 配置auto.offset.resetearliest 历史数据也需要消费查看所有 topicbin/kafka-topics.sh --bootstrap-server localhost:9092 --list删除 topicbin/kafka-topics.sh --bootstrap-server localhost:9092 --delete --topic test-topic执行 --delete 只是打标记不是立刻删磁盘文件或目录test-topic-0、test-topic-1、*-delete* 这些残留。如果存在kafka会不停重启。消费数据存在log.dirscat /opt/kafka/config/server.properties | grep log.dirs # 1. 进入 Kafka 容器 docker exec -it kafka bash # 2. 进入数据目录log.dirs cd /var/lib/kafka/data # 3. 删除坏的 topic 文件夹解决重启根源 rm -rf test-topic-* rm -rf *-delete* # 4. 退出容器 exit生产者发消息bin/kafka-console-producer.sh --topic test-topic --bootstrap-server localhost:9092输入任意内容回车发送。消费者收消息另开一个终端进入容器bin/kafka-console-consumer.sh --topic test-topic --from-beginning --bootstrap-server localhost:9092能收到消息即正常。PHP 使用kafka需要装扩展1安装kafka的扩展之前在安装php-rdkafka之前需要先安装librdkafkagit clone https://github.com/edenhill/librdkafka.git cd librdkafka ./configure make make install2安装rdkafkagit clone https://github.com/arnaud-lb/php-rdkafka.git cd php-rdkafka phpize ./configure --with-php-config/usr/local/php/bin/php-config ###你安装的php下的php-config路径 make make install3、在php.ini中添加行 ###在php下面的/etc/php.ini 编辑extensionrdkafka.so4、重启php后访问php测试页面验证注意事项如果你的服务器上存在多个版本的php编译的时候要将 –with-php-config 指定到目标PHP 版本的安装目录。高级API的特点优点● 高级API写起来简单● 不需要去自行去管理offset系统通过zookeeper自行管理● 不需要管理分区副本等情况系统自动管理● 消费者断线会自动根据上一次记录在 zookeeper中的offset去接着获取数据默认设置5s更新一下 zookeeper 中存的的offset,版本为0.10.2● 可以使用group来区分对访问同一个topic的不同程序访问分离开来不同的group记录不同的offset这样不同程序读取同一个topic才不会因为offset互相影响缺点● 不能自行控制 offset对于某些特殊需求来说● 不能细化控制如分区、副本、zk 等低级API的特点优点● 能够开发者自己控制offset想从哪里读取就从哪里读取。● 自行控制连接分区对分区自定义进行负载均衡● 对 zookeeper 的依赖性降低如offset 不一定非要靠 zk 存储自行存储offset 即可比如存在文件或者内存中缺点● 太过复杂需要自行控制 offset连接哪个分区找到分区 leader 等根据不同的业务需求编写生产者消费者需要选择不同的API接口代码发送消息在基于Kafka的消息中仅仅支持部分简单的类型如StringInteger。但通常使用中需要传递到复杂对象数组队列等可以使用JSON可以很容易的转换为String。?php try { $rcf new RdKafka\Conf(); // 配置groud.id 具有相同 group.id 的consumer 将会处理不同分区的消息所以同一个组内的消费者数量如果订阅了一个topic 那么消费者进程的数量多于这个topic分区的数量是没有意义的。 $rcf-set(group.id, test); $cf new RdKafka\TopicConf(); $cf-set(offset.store.method, broker); //当没有初始偏移量时从哪里开始读取 $cf-set(auto.offset.reset, smallest); $rk new RdKafka\Producer($rcf); $rk-setLogLevel(LOG_DEBUG); $rk-addBrokers(127.0.0.1); $topic $rk-newTopic(ordersMq, $cf); for ($i 0; $i 1000; $i) { $message test kafka message . $i; $topic-produce(0, 0, $message); } } catch (Exception $e) { echo $e-getMessage(); }接收消息?php try { $rcf new RdKafka\Conf(); $rcf-set(group.id, test); $cf new RdKafka\TopicConf(); //$cf-set(offset.store.method, file); $cf-set(auto.offset.reset, smallest); $cf-set(auto.commit.enable, true); $rk new RdKafka\Consumer($rcf); $rk-setLogLevel(LOG_DEBUG); // 指定 broker 地址,多个地址用, 分割 $rk-addBrokers(127.0.0.1); $topic $rk-newTopic(ordersMq, $cf); //$topic-consumeStart(0, RD_KAFKA_OFFSET_BEGINNING); while (true) { $topic-consumeStart(0, RD_KAFKA_OFFSET_STORED); // 第一个参数是分区号 // 第二个参数是超时时间 $msg $topic-consume(0, 1000); // var_dump($msg); if ($msg-err) { echo $msg-errstr(), \n; break; } else { echo $msg-payload, \n; } $topic-consumeStop(0); sleep(1); } } catch (Exception $e) { echo $e-getMessage(); }运行消费者php consumer.phprabbitmq kafka 实际场景选择mq是适合业务中间件 kafka是数据中间件。在实际生产应用中通常会使用kafka作为消息传输的数据管道rabbitmq作为交易数据作为数据传输管道主要的取舍因素则是是否存在丢数据的可能rabbitmq在金融场景中经常使用具有较高的严谨性数据丢失的可能性更小同事具备更高的实时性而kafka优势主要体现在吞吐量上虽然可以通过策略实现数据不丢失但从严谨性角度来讲大不如rabbitmq而且由于kafka保证每条消息最少送达一次有较小的概率会出现数据重复发送的情况Java Spring整合Kafkadependency groupIdorg.springframework.kafka/groupId artifactIdspring-kafka/artifactId version2.7.2/version /dependencyapplication.ymlspring: kafka: bootstrap-servers: localhost:9092 # ✅ 你的 Docker 地址正确 consumer: group-id: test-group # ✅ 消费者组名 auto-offset-reset: earliest # ✅ 从头消费测试非常好用 producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer配置类import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.annotation.EnableKafka; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.core.ProducerFactory; import org.springframework.kafka.listener.ContainerProperties; import java.util.HashMap; import java.util.Map; /** * Kafka的配置类 */ Configuration EnableKafka public class KafkaConfig { Value(${spring.kafka.bootstrap.servers}) private String bootstrapServers; Value(${spring.kafka.group.id}) private String groupId; Value(${spring.kafka.retries}) private String retries; Value(${spring.kafka.concurrency:3}) private Integer concurrency; /** * kafka消息监听器容器的工厂类 * * return */ Bean ConcurrentKafkaListenerContainerFactoryInteger, String kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactoryInteger, String factory new ConcurrentKafkaListenerContainerFactory(); // 3个KafkaMessageListenerContainer并发监听 factory.setConcurrency(concurrency); // 消费者工厂 factory.setConsumerFactory(consumerFactory()); ContainerProperties containerProperties factory.getContainerProperties(); // 当Acknowledgment.acknowledge()方法被调用即提交offset containerProperties.setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); // 调用commitAsync()异步提交 containerProperties.setSyncCommits(false); return factory; } /** * 消费者工厂 * * return */ Bean public ConsumerFactoryInteger, String consumerFactory() { return new DefaultKafkaConsumerFactory(consumerConfigs()); } /** * 消费者拉取消息配置 * * return */ Bean public MapString, Object consumerConfigs() { MapString, Object props new HashMap(16); // kafka集群地址 props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); // groupId props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); // 开启自动提交 props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); // 自动提交offset到zk的时间间隔时间单位是毫秒 // props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 1000); // session超时设置15秒超过这个时间会认为此消费者挂掉将其从消费组中移除 props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 15000); //键的反序列化方式key表示分区 props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); //值的反序列化方式 props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); return props; } /** * 生产者工厂 * * return */ Bean public ProducerFactoryString, String producerFactory() { return new DefaultKafkaProducerFactory(producerConfigs()); } /** * 生产者发送消息配置 * * return */ Bean public MapString, Object producerConfigs() { MapString, Object props new HashMap(); // kafka集群地址 props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); // 消息发送确认方式 props.put(ProducerConfig.ACKS_CONFIG, 1); // 消息发送重试次数 props.put(ProducerConfig.RETRIES_CONFIG, retries); // 重试间隔时间 props.put(ProducerConfig.RETRY_BACKOFF_MS_CONFIG, 1000); // 控制批处理大小单位为字节 props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); // 批量发送延迟为1毫秒启用该功能能有效减少生产者发送消息次数从而提高并发量 props.put(ProducerConfig.LINGER_MS_CONFIG, 1); // 生产者可以使用的总内存字节来缓冲等待发送到服务器的记录 props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432); //键的反序列化方式,key表示分区 props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); //值的反序列化方式 props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); return props; } /** * Kafka模版类用来发送消息 * * return */ Bean public KafkaTemplateString, String kafkaTemplate() { return new KafkaTemplateString, String(producerFactory()); } }消息发布类import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.support.SendResult; import org.springframework.util.concurrent.ListenableFuture; Slf4j public class KafkaProducer { Autowired private KafkaTemplate kafkaTemplate; Value(${spring.kafka.bootstrap.topic}) private String topic; public void sendMessage(String key, String data) { try { log.info(create ag offline,key:{},data: {}, key, data); ListenableFutureSendResultString, String future kafkaTemplate.send(topic, key, data); //future.addCallback(success - log.info(发送消息成功!), failure - log.error(发送消息失败!失败原因是:{}, failure.getMessage())); } catch (Exception e) { log.error(KafkaProducerService sendMessage: {},e:{}, data, e); } } }发消息kafkaProducer.sendMessage(key, JSONObject.toJSONString(kafkaJson, SerializerFeature.WriteMapNullValue));监听消息消费import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.support.Acknowledgment; import org.springframework.stereotype.Component; import java.util.Optional; Component Slf4j public class KafkaConsumer { KafkaListener(topics ${config.kafka.mt_topic}, groupId ${config.kafka.group_id}) public void onMessage(ConsumerRecord?, ? record, Acknowledgment ack) { Optional? message Optional.ofNullable(record.value()); log.info(分区:{},偏移量:{},内容为{}, record.partition(), record.offset(), message); if (message.isPresent()) { try { VosMessage ms VosMessage.parse(message.get()); log.info(消费kafka消息: {}, ms); if (ms ! null) { } } catch (Exception e) { log.error(消费kafka消息 : {} error: {}, message.get(), e); } finally { ack.acknowledge(); } } } }