【发布时间】:2019-10-01 04:47:28
【问题描述】:
我知道必须有办法做到这一点,但我无法弄清楚这一点。读取队列中的所有消息后,我需要停止 kafka 消费者。
有人可以提供这方面的任何信息吗?
【问题讨论】:
标签: apache-kafka kafka-consumer-api
我知道必须有办法做到这一点,但我无法弄清楚这一点。读取队列中的所有消息后,我需要停止 kafka 消费者。
有人可以提供这方面的任何信息吗?
【问题讨论】:
标签: apache-kafka kafka-consumer-api
您可以在启动消费者时传递参数:-consumer-timeout-ms,如果在此期间没有读取任何消息,它将抛出异常。例如,如果在最后 2 秒内没有新消息到达,则停止消费者: kafka.consumer.ConsoleConsumer -consumer-timeout-ms 2000
你可以看到这个和所有其他input options here
【讨论】:
./kafka-console-consumer.sh --bootstrap-server my-host:9092 --new-consumer --topic test-topic --from-beginning --timeout-ms 2000 我试过了,它工作正常,虽然它在退出时会打印一条错误消息和堆栈跟踪。
您可以将 SimpleConsumerShell 与 no-wait-at-logend 选项一起使用。见SystemTools-SimpleConsumerShell
例如:
./kafka-run-class.bat kafka.tools.SimpleConsumerShell --broker-list localhost:9092 --topic kafkademo --partition 0 --no-wait-at-logend
【讨论】:
目前,Kafka 版本 2.11-2.1.1 有一个名为 kafka-console-consumer.sh 的脚本。
它有一个新标志:--timeout-ms。
基本上,这个标志是在没有新日志等待时退出前等待的最长时间。以毫秒为单位。
您可以在阅读所有消息后使用此属性结束您的控制台使用者。
【讨论】:
如果您不打算使用 Scala 客户端,请尝试使用 kafkacat 和 -e 选项告诉它在达到分区结束时退出。
例如消费来自 mytopic 分区 2 的所有消息,然后退出:
$ kafkacat -b mybroker -t mytopic -p 2 -o beginning -e
或者消费最后3000条消息然后退出:
$ kafkacat -b mybroker -t mytopic -p 2 -o -3000 -e
【讨论】: