Perl 6 client for Apache Kafka
Branch: master
Clone or download
mempko Merge pull request #8 from samcv/SPDX-license
Use SPDX identifier in license field of META6.json
Latest commit c3b63cf Apr 25, 2017
Permalink
Type Name Latest commit message Commit time
Failed to load latest commit information.
examples Should be LGPL. Jan 15, 2016
lib/PKafka Fix compilation issue with latest rakudo release. Feb 7, 2017
t Added META6.json test Feb 7, 2017
LICENSE Initial alpha version of the PKafka library. Jan 3, 2016
META6.json
README.md

README.md

PKafka

A Perl 6 client for Apache Kafka. You can use Perl 6's reactive programming features to 'tap' Kafka topics and process messages.

SYNOPSIS

use PKafka::Consumer;
use PKafka::Message;
use PKafka::Producer;

sub MAIN () 
{
    my $brokers = "127.0.0.1";
    my $test = PKafka::Consumer.new( topic=>"test", brokers=>$brokers);
    my $test2 = PKafka::Consumer.new( topic=>"test2", brokers=>$brokers);
    my $producer = PKafka::Producer.new( topic=>"test2", brokers=>$brokers);

    $test.messages.tap(-> $msg 
    {
        given $msg 
        {
            when PKafka::Message
            {
                say "got {$msg.offset}: { $msg.payload-str } ";
                $producer.put("from test '{$msg.payload-str}'");
            }
            when PKafka::EOF
            {
                say "Messages Consumed { $msg.total-consumed}";
            }
            when PKafka::Error
            {
                say "Error {$msg.what}";
                $test.stop;
            }
        }
    });

    $test2.messages.tap(-> $msg 
    {
        given $msg 
        {
            when PKafka::Message
            {
                say "got {$msg.offset}: { $msg.payload-str } ";
            }
        }
    });

    my $t1 = $test.consume-from-beginning(partition=>0);
    my $t2 = $test2.consume-from-beginning(partition=>0);

    await $t1;
    await $t2;
}

DEPENDENCIES

This library wraps librdkafka and it requires it to be installed to function.

WARNING

This library is ALPHA quality software. Please report any bugs and contribute fixes.