-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathcommit-uncommitted.php
More file actions
55 lines (44 loc) · 1.18 KB
/
Copy pathcommit-uncommitted.php
File metadata and controls
55 lines (44 loc) · 1.18 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
<?php
declare(strict_types=1);
require __DIR__ . '/../vendor/autoload.php';
use Amp\TimeoutCancellation;
use Thesis\Kafka\Client;
use Thesis\Kafka\Config;
use Thesis\Kafka\Consumer;
use Thesis\Kafka\Exception;
use function Amp\async;
use function Amp\trapSignal;
/** @var \Psr\Log\LoggerInterface $logger */
$logger = require_once __DIR__ . '/logger.php';
$client = new Client(
new Config(
seeds: [
'kafka-1:9092',
'kafka-2:9092',
'kafka-3:9092',
],
),
$logger,
);
$consumer = $client->createConsumer(['events'], new Consumer\Config(group: new Consumer\GroupConfig(
groupId: 'thesis-consumer',
)));
$future = async(static function () use (
$consumer,
$client,
): void {
trapSignal([\SIGINT, \SIGTERM]);
$consumer->close();
$client->close();
});
try {
while (true) { // @phpstan-ignore while.alwaysTrue
foreach ($consumer->consume(Consumer\ConsumeMode::Poll, new TimeoutCancellation(1.0)) as $batch) {
dump($batch);
}
$consumer->commitUncommitted();
}
} catch (Exception\ConsumerClosed) {
// close() was called from the signal handler above.
}
$future->await();