kafka-php的github地址 https://github.com/jacky5059/kafka-php
生产者produce示例代码
<?php set\_include\_path( implode(PATH_SEPARATOR, array( realpath(\_\_DIR\_\_ . '/../lib'),
get\_include\_path(), ))
); require 'autoloader.php'; $host = 'localhost'; $port = 9092; $topic = 'test'; $producer = new Kafka_Producer($host, $port, Kafka_Encoder::COMPRESSION_NONE); $in = fopen('php://stdin', 'r'); while (true) { echo "\\nEnter comma separated messages:\\n"; $messages = explode(',', fgets($in)); foreach (array_keys($messages) as $k) { //$messages\[$k\] = trim($messages\[$k\]);
} $bytes = $producer->send($messages, $topic); printf("\\nSuccessfully sent %d messages (%d bytes)\\n\\n", count($messages), $bytes);
}
简单消费者simple consumer示例代码
<?php set\_include\_path( implode(PATH_SEPARATOR, array( realpath(\_\_DIR\_\_ . '/../lib'),
get\_include\_path(), ))
); require 'autoloader.php'; $host = 'localhost'; $zkPort = 2181; //zookeeper
$kPort = 9092; //kafka server
$topic = 'test'; $maxSize = 10000000; $socketTimeout = 2; $offset = 0; $partition = 0; $nMessages = 0; $consumer = new Kafka_SimpleConsumer($host, $kPort, $socketTimeout, $maxSize); while (true) { try { //create a fetch request for topic "test", partition 0, current offset and fetch size of 1MB
$fetchRequest = new Kafka_FetchRequest($topic, $partition, $offset, $maxSize); //get the message set from the consumer and print them out
$partialOffset = 0; $messages = $consumer->fetch($fetchRequest); foreach ($messages as $msg) { ++$nMessages; echo "\\nconsumed\[$offset\]\[$partialOffset\]\[msg #{$nMessages}\]: " . $msg->payload(); $partialOffset = $messages->validBytes();
} //advance the offset after consuming each message
$offset += $messages->validBytes(); //echo "\\n---\[Advancing offset to $offset\]------(".date('H:i:s').")";
unset($fetchRequest); //sleep(2);
} catch (Exception $e) { // probably consumed all items in the queue.
echo "\\nERROR: " . get_class($e) . ': ' . $e->getMessage()."\\n".$e->getTraceAsString()."\\n"; sleep(2);
}
}
基于zookeeper的消费者zkconsumer示例代码
<?php set\_include\_path( implode(PATH_SEPARATOR, array( realpath(\_\_DIR\_\_ . '/../lib'),
get\_include\_path(), ))
); require 'autoloader.php'; // zookeeper address (one or more, separated by commas)
$zkaddress = 'localhost:8121'; // kafka topic to consume from
$topic = 'testtopic'; // kafka consumer group
$group = 'testgroup'; // socket buffer size: must be greater than the largest message in the queue
$socketBufferSize = 10485760; //10 MB
// approximate max number of bytes to get in a batch
$maxBatchSize = 20971520; //20 MB
$zookeeper = new Zookeeper($zkaddress); $zkconsumer = new Kafka_ZookeeperConsumer( new Kafka\_Registry\_Topic($zookeeper),
new Kafka\_Registry\_Broker($zookeeper),
new Kafka\_Registry\_Offset($zookeeper, $group),
$topic,
$socketBufferSize ); $messages = array(); try { foreach ($zkconsumer as $message) { // either process each message one by one, or collect them and process them in batches
$messages\[\] = $message; if ($zkconsumer->getReadBytes() >= $maxBatchSize) { break;
}
}
} catch (Kafka\_Exception\_OffsetOutOfRange $exception) { // if we haven't received any messages, resync the offsets for the next time, then bomb out
if ($zkconsumer->getReadBytes() == 0) { $zkconsumer->resyncOffsets(); die($exception->getMessage());
} // if we did receive some messages before the exception, carry on.
} catch (Kafka\_Exception\_Socket_Connection $exception) { // deal with it below
} catch (Kafka_Exception $exception) { // deal with it below
} if (null !== $exception) { // if we haven't received any messages, bomb out
if ($zkconsumer->getReadBytes() == 0) { die($exception->getMessage());
} // otherwise log the error, commit the offsets for the messages read so far and return the data
} // process the data in batches, wait for ACK
$success = doSomethingWithTheMessages($messages); // Once the data is processed successfully, commit the byte offsets.
if ($success) { $zkconsumer->commitOffsets();
} // get an approximate figure on the size of the queue
try { echo "\\nRemaining bytes in queue: " . $consumer->getRemainingSize();
} catch (Kafka\_Exception\_Socket_Connection $exception) { die($exception->getMessage());
} catch (Kafka_Exception $exception) { die($exception->getMessage());
}
Original url: Access
Created at: 2018-10-29 13:04:43
Category: default
Tags: none
未标明原创文章均为采集,版权归作者所有,转载无需和我联系,请注明原出处,南摩阿彌陀佛,知识,不只知道,要得到
java windows火焰图_mob64ca12ec8020的技术博客_51CTO博客 - 在windows下不可行,不知道作者是怎样搞的 监听SpringBoot 服务启动成功事件并打印信息_监听springboot启动完毕-CSDN博客 SpringBoot中就绪探针和存活探针_management.endpoint.health.probes.enabled-CSDN博客 u2u转换板 - 嘉立创EDA开源硬件平台 Spring Boot 项目的轻量级 HTTP 客户端 retrofit 框架,快来试试它!_Java精选-CSDN博客 手把手教你打造一套最牛的知识笔记管理系统! - 知乎 - 想法有重合-理论可参考 安宇雨 闲鱼 机械键盘 客制化 开贴记录 文本 linux 使用find命令查找包含某字符串的文件_beijihukk的博客-CSDN博客_find 查找字符串 ---- mac 也适用 安宇雨 打字音 记录集合 B站 bilibili 自行搭建 开坑 真正的客制化 安宇雨 黑苹果开坑 查找工具包maven pom 引用地 工具网站 Dantelis 介绍的玩轴入坑攻略 --- 关于轴的一些说法 --- 非官方 ---- 心得而已 --- 长期开坑更新 [本人问题][新开坑位]关于自动化测试的工具与平台应用 机械键盘 开团 网站记录 -- 能做一个收集的程序就好了 不过现在没时间 -- 信息大多是在群里发的 - 你要让垃圾佬 都去一个地方看难度也是很大的 精神支柱 [超级前台]sprinbboot maven superdesk-app 记录 [信息有用] [环境准备] [基本完成] [sebp/elk] 给已创建的Docker容器增加新的端口映射 - qq_30599553的博客 - CSDN博客 [正在研究] Elasticsearch, Logstash, Kibana (ELK) Docker image documentation elasticsearch centos 安装记录 及 启动手记 正式服务器 39 elasticsearch 问题合集 不断更新 6.1.1 | 6.5.1 两个版本 博客程序 - 测试 - bug记录 等等问题 laravel的启动过程解析 - lpfuture - 博客园 OAuth2 Server PHP 用 Laravel 搭建带 OAuth2 验证的 RESTful 服务 | Laravel China 社区 - 高品质的 Laravel 和 PHP 开发者社区 利用Laravel 搭建oauth2 API接口 附 Unauthenticated 解决办法 - 煮茶的博客 - SegmentFault 思否 使用 OAuth2-Server-php 搭建 OAuth2 Server - 午时的海 - 博客园 基于PHP构建OAuth 2.0 服务端 认证平台 - Endv - 博客园 Laravel 的 Artisan 命令行工具 Laravel 的文件系统和云存储功能集成 浅谈Chromium中的设计模式--终--Observer模式 浅谈Chromium中的设计模式--二--pre/post和Delegate模式 浅谈Chromium中的设计模式--一--Chromium中模块分层和进程模型 DeepMind 4 Hacking Yourself README.md update 20211011
Laravel China 简书 知乎 博客园 CSDN博客 开源中国 Go Further Ryan是菜鸟 | LNMP技术栈笔记 云栖社区-阿里云 Netflix技术博客 Techie Delight Linkedin技术博客 Dropbox技术博客 Facebook技术博客 淘宝中间件团队 美团技术博客 360技术博客 古巷博客 - 一个专注于分享的不正常博客 软件测试知识传播 - 测试窝 有赞技术团队 阮一峰 语雀 静觅丨崔庆才的个人博客 软件测试从业者综合能力提升 - isTester IBM Java 开发 使用开放 Java 生态系统开发现代应用程序 pengdai 一个强大的博主 HTML5资源教程 | 分享HTML5开发资源和开发教程 蘑菇博客 - 专注于技术分享的博客平台 个人博客-leapMie 流星007 CSDN博客 - 舍其小伙伴 稀土掘金 Go 技术论坛 | Golang / Go 语言中国知识社区
最新评论