RocketMQ integration for Apache Flink. This module includes the RocketMQ source and sink that allows a flink job to either write messages into a topic or read from topics in a flink job.
After topic expanded, Flink jobs cannot consume new queues。Because of follow reasons:
restoredOffsets doesn't have new queues' checkpoint
Flink job start with offset from restoredOffsets wouldn't initialize offset for new queues. And the consumer wouldn't pull messages from new queues because the offsetTable doesn't contain new queues.
After topic expanded, Flink jobs cannot consume new queues。Because of follow reasons: