消费者

August 26, 2021 · View on GitHub

消费者配置

类名:longlang\phpkafka\Consumer\ConsumerConfig

支持构造方法传入数组赋值

配置参数

参数名说明默认值
connectTimeout连接超时时间(单位:秒,支持小数),为-1则不限制-1
sendTimeout发送超时时间(单位:秒,支持小数),为-1则不限制-1
recvTimeout接收超时时间(单位:秒,支持小数),为-1则不限制-1
clientIdKafka 客户端标识,不同的消费者进程请使用不同的设置null
maxWriteAttempts最大写入尝试次数3
client使用哪个 Kafka 客户端类,默认为null时根据场景自动识别null
socket使用哪个 Kafka Socket 类,默认为null时根据场景自动识别null
brokers别名 broker,格式:'127.0.0.1:9092,127.0.0.1:9093'['127.0.0.1:9092','127.0.0.1:9093']null
bootstrapServers别名bootstrapServer,引导服务器,如果配置了该值,会自动连接该服务器,并自动更新 brokers。格式:'127.0.0.1:9092,127.0.0.1:9093'['127.0.0.1:9092','127.0.0.1:9093']null
updateBrokers是否自动更新 brokerstrue
interval未获取消息到消息时,延迟多少秒再次尝试,默认为0则不延迟(单位:秒,支持小数)0
groupId分组 IDnull
memberId用户 IDnull
groupInstanceId分组实例 ID,不同的消费者进程请使用不同的设置null
sessionTimeout如果超时后没有收到心跳信号,则协调器会认为该用户死亡。(单位:秒,支持小数)60
rebalanceTimeout重新平衡组时,协调器等待每个成员重新加入的最长时间(单位:秒,支持小数)。60
topic主题名称,支持同时消费多个主题null
replicaId副本 ID-1
rackId机架编号''
autoCommit自动提交 offsettrue
groupRetry分组操作,匹配预设的错误码时,自动重试次数5
groupRetrySleep分组操作重试延迟,单位:秒1
offsetRetry偏移量操作,匹配预设的错误码时,自动重试次数5
groupHeartbeat分组心跳时间间隔,单位:秒3
autoCreateTopic自动创建主题true
partitionAssignmentStrategy消费者分区分配策略,可选:范围分配-longlang\phpkafka\Consumer\Assignor\RangeAssignor、轮询分配-\longlang\phpkafka\Consumer\Assignor\RoundRobinAssignor、粘性分配-\longlang\phpkafka\Consumer\Assignor\StickyAssignorlonglang\phpkafka\Consumer\Assignor\RangeAssignor
exceptionCallback遇到无法在recv()协程抛出的异常时,调用此回调。格式:function(\Exception $e){}null
minBytes最小字节数1
maxBytes最大字节数128 * 1024 * 1024
maxWait最大等待时间,单位:秒1
saslSASL身份认证信息。为空则不发送身份认证信息 详情[]
sslSSL链接相关信息,为空则不使用SSL 详情null

异步消费(回调)

代码示例:

use longlang\phpkafka\Consumer\ConsumeMessage;
use longlang\phpkafka\Consumer\Consumer;
use longlang\phpkafka\Consumer\ConsumerConfig;

function consume(ConsumeMessage $message)
{
    var_dump($message->getKey() . ':' . $message->getValue());
    // $consumer->ack($message); // autoCommit设为false时,手动提交
}
$config = new ConsumerConfig();
$config->setBroker('127.0.0.1:9092');
$config->setTopic('test'); // 主题名称
$config->setGroupId('testGroup'); // 分组ID
$config->setClientId('test'); // 客户端ID,不同的消费者进程请使用不同的设置
$config->setGroupInstanceId('test'); // 分组实例ID,不同的消费者进程请使用不同的设置
$config->setInterval(0.1);
$consumer = new Consumer($config, 'consume');
$consumer->start();

同步消费

代码示例:

use longlang\phpkafka\Consumer\Consumer;
use longlang\phpkafka\Consumer\ConsumerConfig;

$config = new ConsumerConfig();
$config->setBroker('127.0.0.1:9092');
$config->setTopic('test'); // 主题名称
$config->setGroupId('testGroup'); // 分组ID
$config->setClientId('test_custom'); // 客户端ID,不同的消费者进程请使用不同的设置
$config->setGroupInstanceId('test_custom'); // 分组实例ID,不同的消费者进程请使用不同的设置
$consumer = new Consumer($config);
while(true) {
    $message = $consumer->consume();
    if($message) {
        var_dump($message->getKey() . ':' . $message->getValue());
        $consumer->ack($message); // 手动提交
    }
    sleep(1);
}

SASL支持

相关配置

参数名说明默认值
typeSASL授权对应的类。PLAIN为\longlang\phpkafka\Sasl\PlainSasl::class''
username账号''
password密码''

代码示例:

use longlang\phpkafka\Consumer\Consumer;
use longlang\phpkafka\Consumer\ConsumerConfig;

$config = new ConsumerConfig();
// .... 你的其他配置
$config->setSasl([
    "type"=>\longlang\phpkafka\Sasl\PlainSasl::class,
    "username"=>"admin",
    "password"=>"admin-secret"
]);
$consumer = new Consumer($config);
// ....  你的业务代码

SSL支持

类名:longlang\phpkafka\Config\SslConfig

支持构造方法传入数组赋值

配置参数

参数名说明默认值
open是否开启SSL传输加密false
compression是否开启压缩true
certFilecert证书存放路径''
keyFile私钥存放路径''
passphrasecert证书密码''
peerName服务器主机名。默认为链接的host''
verifyPeer是否校验远端证书false
verifyPeerName是否校验远端服务器名称false
verifyDepth如果证书链条层次太深,超过了本选项的设定值,则终止验证。 默认不校验层级0
allowSelfSigned是否允许自签证书false
cafileCA证书路径''
capathCA证书目录。会自动扫描该路径下所有pem文件''

代码示例:

use longlang\phpkafka\Consumer\Consumer;
use longlang\phpkafka\Consumer\ConsumerConfig;
use longlang\phpkafka\Config\SslConfig;

$config = new ConsumerConfig();
// .... 你的其他配置
$sslConfig = new SslConfig();
$sslConfig->setOpen(true);
$sslConfig->setVerifyPeer(true);
$sslConfig->setAllowSelfSigned(true);
$sslConfig->setCafile("/kafka-client/.github/kafka/cert/ca-cert");
$config->setSsl($sslConfig);
$consumer = new Consumer($config);
// ....  你的业务代码