Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions README-CN.md
Original file line number Diff line number Diff line change
Expand Up @@ -260,6 +260,17 @@ $options->setReconnectPolicy(true,3);
$options->setReconnectPolicy(true,3,100);
```

> Ack 超时时间

* `ack()` 会阻塞等待 broker 返回的 `ACK_RESPONSE`,通过 request ID 进行匹配
* 如果在超时时间内没有收到匹配的响应,`ack()` 会抛出 `RuntimeException`
* 如果在等待响应期间发生重连,重连后重新发送的 ack 同样受此超时时间限制
* 默认值为 30 秒

```php
$options->setAckTimeout(30);
```

> 不循环接收消息,且平滑退出

```php
Expand Down Expand Up @@ -464,6 +475,7 @@ $reader->close();
* setDeadLetterPolicy()
* setSubscriptionInitialPosition()
* setReconnectPolicy()
* setAckTimeout()
* setSchema()
* ReaderOptions
* setTopic()
Expand Down
12 changes: 12 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -272,6 +272,17 @@ $options->setReconnectPolicy(true,3);
$options->setReconnectPolicy(true,3,100);
```

> Ack timeout

* `ack()` blocks waiting for the broker's `ACK_RESPONSE`, matched to the request by request ID
* If no matching response arrives within this timeout, `ack()` throws a `RuntimeException`
* Also applied to each resend after a reconnect occurs while waiting for the response
* Default is 30 seconds

```php
$options->setAckTimeout(30);
```

> Not loop Receive And Smooth exit

```php
Expand Down Expand Up @@ -478,6 +489,7 @@ $reader->close();
* setDeadLetterPolicy()
* setSubscriptionInitialPosition()
* setReconnectPolicy()
* setAckTimeout()
* setSchema()
* ReaderOptions
* setTopic()
Expand Down
138 changes: 112 additions & 26 deletions src/Consumer.php
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,11 @@
use Pulsar\Exception\MessageNotFound;
use Pulsar\Exception\OptionsException;
use Pulsar\Exception\RuntimeException;
use Pulsar\Proto\BaseCommand\Type;
use Pulsar\Proto\CommandAckResponse;
use Pulsar\Proto\CommandMessage;
use Pulsar\Util\Buffer;
use Pulsar\Util\Helper;
use Pulsar\Util\Packer;
use SplPriorityQueue;
use SplQueue;
Expand Down Expand Up @@ -182,13 +186,12 @@ public function receive(bool $loop = true): Message
}
}


// nack
$this->executeInternalNack();

// ping
$this->ping();

if (is_null($response)) {
if (!$loop) {
throw new MessageNotFound();
Expand All @@ -209,28 +212,11 @@ public function receive(bool $loop = true): Message
return $this->receive($loop);
}


$consumer = $this->getPartitionConsumer($commandMessage->getConsumerId());

/**
* @var $messages array<Message>
*/
$messages = Packer::decode($commandMessage, $response->getBuffer(), $consumer->getTopic());

foreach ($messages as $message) {

// Save Options to Message Object
$message->setOptions($this->options);

$this->messageQueue->enqueue($message);
}

$consumer->decrement(sizeof($messages));
$this->enqueueCommandMessage($commandMessage, $response->getBuffer());

return $this->messageQueue->dequeue();
}


/**
* @return array<Message>
* @throws IOException
Expand All @@ -248,18 +234,87 @@ public function batchReceive(bool $loop = true): array
}

/**
* Sends the CommandAck and blocks for its CommandAckResponse, correlated by request ID.
*
* While waiting, any MESSAGE frame that arrives first (the broker may deliver one before
* the ACK_RESPONSE, since the connection is asynchronous) is queued rather than discarded.
*
* Unlike receive(), this does not reconnect on a dropped connection: a lost connection
* mid-ack throws IOException immediately, regardless of ConsumerOptions::getReconnectPolicy().
*
* This is deliberate, not an oversight. A dropped connection during ack() means the
* broker may already have decided this consumer is gone and redelivered the message to
* another consumer on the same subscription (Shared/Key_Shared). Throwing immediately
* gives the caller an honest, timely signal instead of a client library quietly retrying
* underneath it.
*
* @param Message $message
* @return void
* @throws \Exception
* @return CommandAckResponse|null
* @throws IOException
* @throws RuntimeException
*/
public function ack(Message $message)
public function ack(Message $message): ?CommandAckResponse
{
if (!$message->canAck()) {
return;
return null;
}

// send CommandAck
$this->getPartitionConsumer($message->getConsumerID())->ack($message);
$requestId = Helper::getRequestID();
$this->getPartitionConsumer($message->getConsumerID())->ack($message, $requestId);

$deadline = microtime(true) + $this->options->getAckTimeout();

do {
$remaining = $deadline - microtime(true);
if ($remaining <= 0) {
throw new RuntimeException('Timed out waiting for ACK response.');
}

$response = $this->eventloop->wait((int) ceil($remaining));

if (null === $response) {
continue;
}

$baseCommand = $response->getBaseCommand();

$commandType = $baseCommand->getType();

if (Type::CLOSE_CONSUMER_VALUE === $commandType->value()) {
// only abort if it's the consumer this message belongs to; the connection
// may be shared with other partition consumers that closed independently
if ($baseCommand->getCloseConsumer()->getConsumerId() === $message->getConsumerID()) {
throw new RuntimeException(
'The consumer was closed before the message acknowledgment was confirmed.'
);
}
continue;
}

if (Type::MESSAGE_VALUE === $commandType->value()) {
$this->enqueueCommandMessage($baseCommand->getMessage(), $response->getBuffer());
continue;
}

if (Type::ACK_RESPONSE_VALUE === $commandType->value()) {
$ackResponse = $baseCommand->getAckResponse();

if ($ackResponse->getRequestId() !== $requestId) {
throw new RuntimeException('ACK response request ID does not match.');
}

if ($ackResponse->hasError()) {
$msg = $ackResponse->hasMessage() ? $ackResponse->getMessage() : $ackResponse->getError()->name();
throw new RuntimeException(
sprintf('The broker rejected the acknowledgment: %s', $msg),
$ackResponse->getError()->value()
);
}

return $ackResponse;
}

} while (true);
}


Expand Down Expand Up @@ -358,4 +413,35 @@ protected function getPartitionConsumer(int $consumerID): PartitionConsumer
return $this->consumers[ $consumerID ];
}

/**
* Decodes a MESSAGE frame's payload into Message objects, queues them locally, and
* decrements the partition's available flow-control permits accordingly.
*
* Shared by receive() and ack()'s wait loop, since a MESSAGE frame can arrive while
* ack() is waiting on the same connection for an unrelated ACK_RESPONSE -- it must be
* queued here rather than discarded, or the message would be silently lost.
*
* @param CommandMessage $commandMessage
* @param Buffer $buffer
* @return void
*/
private function enqueueCommandMessage(CommandMessage $commandMessage, Buffer $buffer)
{
$consumer = $this->getPartitionConsumer($commandMessage->getConsumerId());

/**
* @var array<Message> $messages
*/
$messages = Packer::decode($commandMessage, $buffer, $consumer->getTopic());

foreach ($messages as $message) {
// Save Options to Message Object
$message->setOptions($this->options);

$this->messageQueue->enqueue($message);
}

$consumer->decrement(sizeof($messages));
}

}
27 changes: 27 additions & 0 deletions src/ConsumerOptions.php
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,13 @@ final class ConsumerOptions extends Options
*/
const ENABLE_RECONNECT = 'reconnect';

/**
* Max seconds to wait for an ACK_RESPONSE before giving up, default 30 seconds
*
* @var int
*/
const ACK_TIMEOUT = 'ack_timeout';

/**
* @param array $topics
* @return void
Expand Down Expand Up @@ -259,6 +266,26 @@ public function getNackRedeliveryDelay()
}


/**
* @param int $seconds
* @return void
*/
public function setAckTimeout(int $seconds)
{
$this->data[ self::ACK_TIMEOUT ] = $seconds;
}


/**
* @return int|mixed
*/
public function getAckTimeout()
{
$timeout = $this->data[ self::ACK_TIMEOUT ] ?? 30;
return $timeout <= 0 ? 30 : $timeout;
}


/**
* @return mixed|string
* @throws \Exception
Expand Down
8 changes: 6 additions & 2 deletions src/PartitionConsumer.php
Original file line number Diff line number Diff line change
Expand Up @@ -12,10 +12,12 @@


use Pulsar\Exception\IOException;
use Pulsar\Exception\RuntimeException;
use Pulsar\IO\AbstractIO;
use Pulsar\Proto\BaseCommand\Type;
use Pulsar\Proto\CommandAck;
use Pulsar\Proto\CommandAck\AckType;
use Pulsar\Proto\CommandAckResponse;
use Pulsar\Proto\CommandCloseConsumer;
use Pulsar\Proto\CommandFlow;
use Pulsar\Proto\CommandRedeliverUnacknowledgedMessages;
Expand Down Expand Up @@ -155,10 +157,11 @@ public function flow()

/**
* @param Message $message
* @return void
* @param int $requestId
* @return CommandAckResponse
* @throws IOException
*/
public function ack(Message $message)
public function ack(Message $message, int $requestId)
{
// send CommandAck
$command = new CommandAck();
Expand All @@ -167,6 +170,7 @@ public function ack(Message $message)
$command->addMessageId($message->getMessageIdData());
$command->setTxnidLeastBits(null);
$command->setTxnidMostBits(null);
$command->setRequestId($requestId);
$this->connection->writeCommand(Type::ACK(), $command);
}

Expand Down