新闻详情

新闻详情

首页 / 资讯中心 / 详情

【基于 Swoole+Hyperf 的微服务实战】第七周·周三:Kafka 与高吞吐消息场景

发布时间:2026/9/25 18:32:39来源:尧图网络
【基于 Swoole+Hyperf 的微服务实战】第七周·周三:Kafka 与高吞吐消息场景
【基于 SwooleHyperf 的微服务实战】第七周·周三Kafka 与高吞吐消息场景今天我们进入第七周周三主题是Kafka 与高吞吐消息场景。前两天我们基于 RabbitMQ 构建了可靠的消息通信但 RabbitMQ 在处理海量日志、流式数据时可能会遇到吞吐瓶颈。今天我们将引入Apache Kafka——一个分布式、高吞吐、支持持久化的流处理平台。我们将使用hyperf/kafka组件实现高并发日志收集将用户行为日志写入 Kafka并由消费者实时处理体验 Kafka 在海量数据下的强劲性能。今日目标理解 Kafka 的核心概念Topic、Partition、Broker、Consumer Group、Offset。使用 Docker 部署 KafkaKRaft 模式无需 ZooKeeper。安装hyperf/kafka配置生产者和消费者。实现一个用户行为日志收集场景HTTP 请求日志通过 Kafka 异步发送消费者批量写入数据库或对象存储。通过压测对比 Kafka 与 RabbitMQ 的吞吐能力感受 Kafka 在流式数据处理中的优势。一、环境准备部署 Kafka约 30 分钟1. 添加 Kafka 服务到 Docker ComposeKafka 3.x 支持KRaft模式无 ZooKeeper。编辑swoole-course/docker-compose.yml增加kafka:image:confluentinc/cp-kafka:7.5.0container_name:kafka-labports:-9092:9092-9093:9093# 内部监听可忽略environment:KAFKA_NODE_ID:1KAFKA_LISTENER_SECURITY_PROTOCOL_MAP:CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXTKAFKA_LISTENERS:PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093KAFKA_ADVERTISED_LISTENERS:PLAINTEXT://kafka:9092KAFKA_PROCESS_ROLES:broker,controllerKAFKA_CONTROLLER_QUORUM_VOTERS:1kafka:9093KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR:1KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR:1KAFKA_TRANSACTION_STATE_LOG_MIN_ISR:1restart:unless-stopped启动 Kafkadocker-composeup-dkafka确认 Kafka 启动成功可能需要几十秒dockerlogs kafka-lab|grepstarted2. 安装 hyperf/kafka 组件进入 PHP 容器docker-composeexecswoolebashcd/var/www/hyperf-appcomposerrequire hyperf/kafka发布配置php bin/hyperf.php vendor:publish hyperf/kafka这将在config/autoload/kafka.php生成默认配置。二、知识核心Kafka 架构与高吞吐原理约 1 小时1. Kafka 基本概念Topic消息分类类似 RabbitMQ 的 routing key 和队列的组合。每个 Topic 可以有多个分区。Partition一个 Topic 被分割为多个分区每个分区是一个有序的、不可变的消息序列。分区是 Kafka 并行处理的核心可以分布在不同的 Broker 上。BrokerKafka 服务节点。一个集群由多个 Broker 组成每个 Broker 承担一部分分区。Producer发布消息到指定 Topic可以指定分区键或由 Kafka 自动分区。Consumer Group消费者组组内多个消费者协同消费一个 Topic 的多个分区每个分区只能被组内一个消费者消费。水平扩展消费者数量时分区数要与之匹配。Offset消费者在分区内的位置标记用于记录已消费位置支持回溯和重放。2. 与 RabbitMQ 的对比特性RabbitMQKafka设计目标可靠消息传递灵活路由高吞吐、流式数据处理消息模型交换机队列支持复杂路由发布-订阅Topic/Partition吞吐量万级 QPS百万级 QPS持久化内存磁盘默认可靠默认磁盘顺序写入极快消息回溯不支持消费后删除支持可按 Offset 重放历史数据事务支持支持幂等生产者事务适用场景业务命令、RPC 回调、延迟消息日志收集、实时流计算、事件溯源今天我们的场景是用户行为日志收集非常适合 Kafka因为日志量大、允许少量延迟、需要持久化和顺序保证。3. hyperf/kafka 组件基于longlang/phpkafkaSwoole 协程化的 Kafka 客户端提供协程化的生产者支持批量发送。消费者可以按分区消费并手动提交 Offset。支持 SASL 认证、SSL 加密等。三、实战构建用户行为日志 Kafka 管道约 2.5 小时步骤 1配置 Kafka 连接编辑config/autoload/kafka.php?phpreturn[default[hostenv(KAFKA_HOST,kafka),port(int)env(KAFKA_PORT,9092),pool[min_connections1,max_connections10,connect_timeout10.0,wait_timeout3.0,heartbeat-1,max_idle_time60,],],];步骤 2创建日志生产者新建app/Kafka/Producer/LogProducer.php?phpnamespaceApp\Kafka\Producer;useHyperf\Kafka\AbstractProducer;useHyperf\Kafka\Annotation\Producer;#[Producer(topic:user_behavior_logs,name:LogProducer)]classLogProducerextendsAbstractProducer{// 无需额外代码基类已提供 send 方法}说明AbstractProducer提供了send(string $payload, string $key null)方法$key用于分区路由相同 key 进入同一分区保证顺序。步骤 3在业务中发送日志消息我们可以在网关的 JWT 中间件或控制器中记录用户每次请求的行为日志。为了演示在hyperf-app的ArticleController::show()方法中发送一条日志消息。修改app/Controller/ArticleController.phpuseApp\Kafka\Producer\LogProducer;useHyperf\Di\Annotation\Inject;#[Inject]privateLogProducer$logProducer;publicfunctionshow(int$id){// 原有业务逻辑...$articleArticle::with([category,tags])-find($id);if(!$article){/* ... */}// 发送用户行为日志到 Kafka$userId$this-request-getHeaderLine(X-User-Id)??0;$logDatajson_encode([user_id(int)$userId,actionview_article,article_id$id,timestamptime(),ip$this-request-getServerParams()[remote_addr]??,]);$this-logProducer-send($logData,user_.$userId);// 以用户ID为分区键同一用户日志有序// 返回文章...}这样每次查看文章详情时一条日志就会发送到 Kafka Topicuser_behavior_logs。步骤 4创建日志消费者新建app/Kafka/Consumer/LogConsumer.php?phpnamespaceApp\Kafka\Consumer;useHyperf\Kafka\AbstractConsumer;useHyperf\Kafka\Annotation\Consumer;useHyperf\Kafka\Result;uselonglang\phpkafka\Consumer\ConsumeMessage;#[Consumer(topic:user_behavior_logs,groupId:log-processor,name:LogConsumer,nums:2)]classLogConsumerextendsAbstractConsumer{publicfunctionconsume(ConsumeMessage$message):string{$datajson_decode($message-getValue(),true);if(!$data){returnResult::ACK;}$userId$data[user_id]??0;$action$data[action]??;echo[Kafka消费者] 处理日志: 用户{$userId}{$action}文章{$data[article_id]}\n;// 模拟批量写入数据库或 ES这里仅打印// 实际可写入 MySQL 日志表或发送到 Elasticsearch 用于分析returnResult::ACK;}}要点groupId消费者组名称同一组的消费者分配不同分区保证每条消息只被组内一个消费者处理。nums启动的消费者协程数应小于等于 Topic 的分区数才能充分利用并行性。ConsumeMessage提供getValue()、getKey()、getOffset()等方法。步骤 5配置消费者进程编辑config/autoload/processes.php添加 Kafka 消费者进程如果没有则创建文件?phpreturn[\Hyperf\Kafka\Process\ConsumerProcess::class,];重启hyperf-appphp bin/hyperf.php start消费者进程会启动自动加入log-processor消费者组消费user_behavior_logs分区。步骤 6创建 Topic手动或自动Kafka 默认允许自动创建 Topicauto.create.topics.enabletrue但建议手动创建以控制分区数。在 Kafka 容器中执行dockerexec-itkafka-lab kafka-topics--create--topicuser_behavior_logs --bootstrap-server kafka:9092--partitions3--replication-factor1查看 Topicdockerexec-itkafka-lab kafka-topics--list--bootstrap-server kafka:9092四、成果测试与高吞吐验证约 1 小时1. 功能验证请求文章详情接口curlhttp://localhost:9501/articles/1观察控制台消费者日志输出确认日志消息被消费。查看 Kafka 消息使用控制台消费者dockerexec-itkafka-lab kafka-console-consumer--topicuser_behavior_logs --bootstrap-server kafka:9092 --from-beginning能看到历史消息。2. 压测吞吐量对比使用wrk或ab对文章详情接口进行高并发压测同时观察 Kafka 生产者和消费者的处理速度。# 压测 20000 请求并发 100wrk-t4-c100-d60shttp://localhost:9501/articles/1RabbitMQ场景下昨天消息吞吐受单队列限制压测时可能堆积。Kafka场景下多分区并行消费者组中多个协程并发消费吞吐极高。可以在 Kafka 管理工具如kafka-consumer-groups查看消费延迟dockerexec-itkafka-lab kafka-consumer-groups --bootstrap-server kafka:9092--grouplog-processor--describe观察LAG未消费消息数是否迅速降为零。3. 测试消费者组和分区重平衡启动两个hyperf-app实例不同 Worker 或多个容器它们属于同一消费者组观察 Kafka 自动将分区均匀分配。停止其中一个实例剩余实例会自动接管全部分区Rebalance。4. 测试清单检验项方法通过标准生产者发送消息查看 Kafka 控制台消费者或 Topic 消息数消息成功进入 Topic消费者接收消息查看应用日志输出打印出日志信息且LAG近 0分区键路由相同user_id的消息Kafka 工具查看分区进入同一分区消费者组并行高并发下消费速度跟得上生产速度无明显堆积Offset 自动提交重启消费者后不会重复消费已处理的消息消息数量不重复需手动确认配置高吞吐稳定性长时间压测系统无 OOM连接正常持续稳定常见问题Topic 不存在生产者发送时会自动创建如果开启但建议手动创建指定分区数。消费者不消费检查 Group ID 是否一致Topic 名称拼写以及分区数是否大于消费者nums。连接 Kafka 失败确认kafka主机名可达PHP 容器内 ping kafka或使用host.docker.internal。五、今日作业与学习产出提交代码将LogProducer、LogConsumer、Kafka 配置文件、进程配置等提交。完善日志场景实现一个日志持久化消费者将行为日志批量写入 MySQL使用事务或写入 Elasticsearch 供 Kibana 可视化。添加异常行为检测消费者分析日志若某用户 1 分钟内请求超过 100 次发送告警消息到 RabbitMQ 或直接钉钉。学习笔记画出 Kafka 的分区、消费者组模型图说明如何水平扩展消费者。对比 RabbitMQ 与 Kafka 的使用场景明确在架构中何时选择哪一个。挑战任务使用Kafka Streams或KSQL实现对日志的实时聚合统计如热门文章 TopN并暴露 API 供查询。配置Kafka Connect将日志直接写入 S3 对象存储实现数据湖入湖。通过今天的学习你已掌握高吞吐场景下的异步消息利器 Kafka。现在你的微服务同时拥有了 RabbitMQ可靠命令和 Kafka流式日志两种消息引擎可以根据业务特点灵活选择。明天我们将结合事件驱动设计实现内部服务的完全解耦与最终一致性。
网站建设高端定制企业官网
RELATED

相关资讯

更多精彩内容,欢迎继续阅读

较早相关资讯

最新相关资讯

实时Linux为什么要隔离内核?从系统调用到内核抢占,看懂硬实时系统的最后一道实时性边界 2026/9/25 19:09:32

实时Linux为什么要隔离内核?从系统调用到内核抢占,看懂硬实时系统的最后一道实时性边界

前面的文章我们一直在讨论一个核心问题:如何让实时任务不被其他任务干扰?从CPU核心隔离,到IRQ隔离;从内存、Cache、DMA等资源控制,到实时IPC;再到Watchdog和故障隔离,我们实际上已经逐渐建立起了…

阅读更多 →
实时Linux中的Watchdog到底有什么用?从任务超时到系统自恢复,看懂硬实时系统如何处理“失控任务” 2026/9/25 19:09:26

实时Linux中的Watchdog到底有什么用?从任务超时到系统自恢复,看懂硬实时系统如何处理“失控任务”

在工业机器人、智能制造、无人系统、飞控、能源控制等场景中,实时操作系统面对的并不只是一个问题:“任务能不能按时运行?”还有一个更加现实的问题:“如果任务没有按时运行,系统怎么办?”例如,…

阅读更多 →
用Trae把原型图变成可执行HTML:TaoToken统一Key接入AI编程工作流 2026/9/25 19:09:26

用Trae把原型图变成可执行HTML:TaoToken统一Key接入AI编程工作流

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
多Agent协作架构与任务调度实战:从单Agent到复杂AI协同系统 2026/9/25 19:09:00

多Agent协作架构与任务调度实战:从单Agent到复杂AI协同系统

1. 多Agent协作到底在解决什么问题单Agent跑任务,跑到一定复杂度就会撞墙。我最早做自动化流程的时候,一个Agent包揽需求解析、资料检索、代码生成、结果校验,提示词写到三千字,工具挂了十几个,结果就是:它…

阅读更多 →
SKILL编排:给存量代码做最小侵入的微创手术 2026/9/25 19:09:00

SKILL编排:给存量代码做最小侵入的微创手术

1. “散装 AI”,正在成为团队里最贵的隐性负债如果你和我一样,在过去一年里反复安慰自己“AI 至少能帮我们写点单测、解释几段历史代码”,那你大概率也注意到了另一件事:团队里的 AI 能力,正在以极其混乱的方式野蛮生长…

阅读更多 →
图像去雾数据集选型与工业落地指南 2026/9/25 19:08:53

图像去雾数据集选型与工业落地指南

1. 项目概述:为什么“图像去雾数据集总汇”不是一张表格,而是一套工程能力图像去雾数据集总汇——这名字听起来像一份静态清单,但在我过去八年做计算机视觉项目落地的过程中,它从来就不是简单罗列几个链接、几行下载命令的事。它本…

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

联系尧图顾问,获取一对一建站咨询

立即免费咨询 📞 400-888-8888
📞 ✉