1. Go to this page and download the library: Download micromus/kafka-bus-commiter library. Choose the download type require.
2. Extract the ZIP file and open the index.php.
3. Add this code to the index.php.
<?php
require_once('vendor/autoload.php');
/* Start to develop here. Best regards https://php-download.com/ */
micromus / kafka-bus-commiter example snippets
use Micromus\KafkaBus\Bus;
use Micromus\KafkaBus\Consumers\Router\ConsumerRoutesBuilder;
use Micromus\KafkaBus\Consumers\Router\RouteInfo;
use Micromus\KafkaBus\Topics\Topic;
use Micromus\KafkaBus\Topics\TopicRegistry;
use Micromus\KafkaBusCommiter\Middleware\ConsumerCommiterMiddleware;
$topicRegistry = (new TopicRegistry())
->add(new Topic('production.fact.products.1', 'products'));
$consumerRoutes = ConsumerRoutesBuilder::make($topicRegistry)
->add(new RouteInfo('products', new YourMessageHandler()))
->build();
$workerRegistry = (new Bus\Listeners\Workers\MemoryWorkerRegistry())
->add(
new Bus\Listeners\Workers\Worker(
name: 'default-listener',
routes: $consumerRoutes,
options: new Bus\Listeners\Workers\Options(
middleware: [
new ConsumerCommiterMiddleware(new NativeMessageRepository(new DatabaseRepositorySource))
]
)
)
);
new ConsumerCommiterMiddleware(
repository: $repository, // ogger
maxAttempt: 3, // max attempts, -1 = unlimited
)
use Micromus\KafkaBusCommiter\Attempt;
use Micromus\KafkaBusCommiter\Interfaces\RepositorySourceInterface;
class DatabaseRepositorySource implements RepositorySourceInterface
{
/**
* Returns the current attempt for a given key.
*/
public function get(string $key): ?Attempt
{
// ...
}
/**
* Increments the number of failed read attempts for a key.
*/
public function increment(string $key): void
{
// ...
}
/**
* Marks the message key as successfully processed.
*/
public function commit(string $key): void
{
// ...
}
}
use Micromus\KafkaBus\Interfaces\Producers\Messages\ProducerMessageInterface;
use Micromus\KafkaBusCommiter\Interfaces\HasIdempotency;
final readonly class ProductCreated implements ProducerMessageInterface, HasIdempotency
{
public function __construct(
private string $productId,
private string $payload,
) {
}
public function toPayload(): string
{
return $this->payload;
}
public function getIdempotencyKey(): string
{
return $this->productId;
}
}
use Micromus\KafkaBus\Bus\Publishers\Router\Options;
use Micromus\KafkaBus\Bus\Publishers\Router\PublisherRoutesBuilder;
use Micromus\KafkaBusCommiter\Middleware\ProducerIdempotencyMiddleware;
$publisherRoutes = PublisherRoutesBuilder::make($topicRegistry)
->add(
ProductCreated::class,
'products',
new Options(middleware: [new ProducerIdempotencyMiddleware()])
)
->build();
use Micromus\KafkaBusCommiter\Middleware\ConsumerCommiterMiddleware;
use Micromus\KafkaBusCommiter\Repositories\IdempotencyMessageRepository;
$repository = new IdempotencyMessageRepository(new DatabaseRepositorySource());
new ConsumerCommiterMiddleware($repository, maxAttempt: 3);
Loading please wait ...
Before you can download the PHP files, the dependencies should be resolved. This can take some minutes. Please be patient.