File tree Expand file tree Collapse file tree 1 file changed +3
-3
lines changed
Expand file tree Collapse file tree 1 file changed +3
-3
lines changed Original file line number Diff line number Diff line change 66use Illuminate \Queue \Connectors \ConnectorInterface ;
77use Rapide \LaravelQueueKafka \Queue \KafkaQueue ;
88use RdKafka \Conf ;
9- use RdKafka \Consumer ;
9+ use RdKafka \KafkaConsumer ;
1010use RdKafka \Producer ;
1111use RdKafka \TopicConf ;
1212
@@ -42,7 +42,7 @@ public function connect(array $config)
4242
4343 /** @var TopicConf $topicConf */
4444 $ topicConf = $ this ->container ->makeWith ('queue.kafka.topic_conf ' , []);
45- $ topicConf ->set ('auto.offset.reset ' , 'smallest ' );
45+ $ topicConf ->set ('auto.offset.reset ' , 'largest ' );
4646
4747 /** @var Conf $conf */
4848 $ conf = $ this ->container ->makeWith ('queue.kafka.conf ' , []);
@@ -52,7 +52,7 @@ public function connect(array $config)
5252 $ conf ->set ('offset.store.method ' , 'broker ' );
5353 $ conf ->setDefaultTopicConf ($ topicConf );
5454
55- /** @var Consumer $consumer */
55+ /** @var KafkaConsumer $consumer */
5656 $ consumer = $ this ->container ->makeWith ('queue.kafka.consumer ' , ['conf ' => $ conf ]);
5757
5858 return new KafkaQueue (
You can’t perform that action at this time.
0 commit comments