Movatterモバイル変換


[0]ホーム

URL:


Skip to content

Navigation Menu

Sign in
Appearance settings

Search code, repositories, users, issues, pull requests...

Provide feedback

We read every piece of feedback, and take your input very seriously.

Saved searches

Use saved searches to filter your results more quickly

Sign up
Appearance settings
This repository was archived by the owner on Jun 10, 2022. It is now read-only.
/kafka-phpPublic archive

Commit8ddac65

Browse files
committed
fixed#104 and add delayQueue producer example
1 parent873e6f1 commit8ddac65

File tree

4 files changed

+67
-6
lines changed

4 files changed

+67
-6
lines changed

‎example/ControlProducer.php‎

Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
1+
<?php
2+
require'../vendor/autoload.php';
3+
date_default_timezone_set('PRC');
4+
useMonolog\Logger;
5+
useMonolog\Handler\StdoutHandler;
6+
// Create the logger
7+
$logger =newLogger('my_logger');
8+
// Now add some handlers
9+
$logger->pushHandler(newStdoutHandler());
10+
11+
$config = \Kafka\ProducerConfig::getInstance();
12+
$config->setMetadataRefreshIntervalMs(1000);
13+
$config->setMetadataBrokerList('10.13.4.162:9192');
14+
$config->setBrokerVersion('0.9.0.1');
15+
$config->setRequiredAck(1);
16+
$config->setIsAsyn(true);
17+
$config->setProduceInterval(500);
18+
19+
class Message {
20+
private$message;
21+
22+
publicfunctiongetMessage() {
23+
return$this->message;
24+
}
25+
publicfunctionsetMessage($message) {
26+
$this->message =$message;
27+
}
28+
}
29+
30+
// control message send interval time
31+
$message =newMessage;
32+
\Amp\repeat(function ()use ($message){
33+
$message->setMessage(array(
34+
array(
35+
'topic' =>'test',
36+
'value' =>'test....message.' .time(),
37+
'key' =>'',
38+
),
39+
));
40+
},3000);
41+
42+
$producer =new \Kafka\Producer(function()use($message) {
43+
$tmp =$message->getMessage();
44+
$message->setMessage([]);
45+
return$tmp;
46+
});
47+
$producer->setLogger($logger);
48+
$producer->success(function($result) {
49+
var_dump($result);
50+
});
51+
$producer->error(function($errorCode,$context) {
52+
var_dump($errorCode);
53+
});
54+
$producer->send();

‎src/Kafka/Consumer/State.php‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -178,6 +178,10 @@ public function succRun($key, $context = null)
178178
break;
179179
caseself::REQUEST_OFFSET:
180180
caseself::REQUEST_FETCH:
181+
if (!isset($this->callStatus[$key]['context'])) {
182+
$this->callStatus[$key]['status'] = (self::STATUS_LOOP |self::STATUS_FINISH);
183+
break;
184+
}
181185
unset($this->callStatus[$key]['context'][$context]);
182186
$contextStatus =$this->callStatus[$key]['context'];
183187
if (empty($contextStatus)) {

‎src/Kafka/Producer/Process.php‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -88,7 +88,8 @@ public function init()
8888
));
8989
$this->state->init();
9090

91-
if (!empty($broker->getTopics())) {
91+
$topics =$broker->getTopics();
92+
if (!empty($topics)) {
9293
$this->state->succRun(\Kafka\Producer\State::REQUEST_METADATA);
9394
}
9495
}
@@ -113,6 +114,7 @@ public function start()
113114
if ($this->error) {
114115
call_user_func($this->error,1000);
115116
}
117+
\Amp\cancel($watcherId);
116118
\Amp\stop();
117119
},$msInterval =$config->getRequestTimeout());
118120
};

‎src/Kafka/Protocol/Fetch.php‎

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -272,11 +272,12 @@ protected function decodeMessage($data, $messageSize)
272272
$attr =self::unpack(self::BIT_B8,substr($data,$offset,1));
273273
$offset +=1;
274274
$timestamp =0;
275-
$version =$this->getApiVersion(self::FETCH_REQUEST);
276-
if ($version ==self::API_VERSION2) {
277-
$timestamp =self::unpack(self::BIT_B64,substr($data,$offset,8));
278-
$offset +=8;
279-
}
275+
//$version = $this->getApiVersion(self::FETCH_REQUEST);
276+
//if ($version == self::API_VERSION2) {
277+
// $timestamp = self::unpack(self::BIT_B64, substr($data, $offset, 8));
278+
// $offset += 8;
279+
//}
280+
280281
$keyRet =$this->decodeString(substr($data,$offset),self::BIT_B32);
281282
$offset +=$keyRet['length'];
282283
$valueRet =$this->decodeString(substr($data,$offset),self::BIT_B32);

0 commit comments

Comments
 (0)

[8]ページ先頭

©2009-2025 Movatter.jp