Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
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
Next Next commit
Pass whole message for deserialization
  • Loading branch information
TiMESPLiNTER committed Sep 26, 2019
commit 7e7e1cd5a7bf595ec97c9d3b6691f6eb09ce7e1c
6 changes: 4 additions & 2 deletions pkg/rdkafka/JsonSerializer.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,8 @@

namespace Enqueue\RdKafka;

use RdKafka\Message as VendorMessage;

class JsonSerializer implements Serializer
{
public function toString(RdKafkaMessage $message): string
Expand All @@ -25,9 +27,9 @@ public function toString(RdKafkaMessage $message): string
return $json;
}

public function toMessage(string $string): RdKafkaMessage
public function toMessage(VendorMessage $message): RdKafkaMessage
{
$data = json_decode($string, true);
$data = json_decode($message->payload, true);
if (JSON_ERROR_NONE !== json_last_error()) {
throw new \InvalidArgumentException(sprintf(
'The malformed json given. Error %s and message %s',
Expand Down
2 changes: 1 addition & 1 deletion pkg/rdkafka/RdKafkaConsumer.php
Original file line number Diff line number Diff line change
Expand Up @@ -164,7 +164,7 @@ private function doReceive(int $timeout): ?RdKafkaMessage
case RD_KAFKA_RESP_ERR__TIMED_OUT:
break;
case RD_KAFKA_RESP_ERR_NO_ERROR:
$message = $this->serializer->toMessage($kafkaMessage->payload);
$message = $this->serializer->toMessage($kafkaMessage);
$message->setKey($kafkaMessage->key);
$message->setPartition($kafkaMessage->partition);
$message->setKafkaMessage($kafkaMessage);
Expand Down
4 changes: 3 additions & 1 deletion pkg/rdkafka/Serializer.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,11 @@

namespace Enqueue\RdKafka;

use RdKafka\Message as VendorMessage;

interface Serializer
{
public function toString(RdKafkaMessage $message): string;

public function toMessage(string $string): RdKafkaMessage;
public function toMessage(VendorMessage $string): RdKafkaMessage;
}