Использование RabbitMQ в качестве источника данных Flink DataStream без автоматического создания очереди RabbitMQ - PullRequest
0 голосов
/ 13 мая 2018

Когда я использую RabbitMQ в качестве источника данных Flink DataStream, как сказано в документации Flink.

final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// checkpointing is required for exactly-once or at-least-once guarantees
env.enableCheckpointing(...);

final RMQConnectionConfig connectionConfig = new RMQConnectionConfig.Builder()
    .setHost("localhost")
    .setPort(5000)
    ...
    .build();

final DataStream<String> stream = env
    .addSource(new RMQSource<String>(
    connectionConfig,            // config for the RabbitMQ connection
    "queueName",                 // name of the RabbitMQ queue to consume
    true,                        // use correlation ids; can be false if only at-least-once is required
    new SimpleStringSchema()))   // deserialization schema to turn messages into Java objects
.setParallelism(1);              // non-parallel source is only required for exactly-once

Этот код подключится к RabbitMQ и автоматически создаст очередь "queueName". Поэтому у меня возникла проблема.Очередь RabbitMQ уже существует, я создал ее раньше.Я не хочу, чтобы Флинк попытался создать снова.И проблема в том, что Flink создает очередь без каких-либо параметров, это конфликтует с очередью, которую я создал ранее.Вот исключение:

    Caused by: com.rabbitmq.client.ShutdownSignalException: channel error; protocol method: #method<channel.close>(reply-code=406, reply-text=PRECONDITION_FAILED - inequivalent arg 'x-message-ttl' for queue 'queueName' in vhost '/': received none but current is the value '604800000' of type 'long', class-id=50, method-id=10)
at com.rabbitmq.utility.ValueOrException.getValue(ValueOrException.java:66)
at com.rabbitmq.utility.BlockingValueOrException.uninterruptibleGetValue(BlockingValueOrException.java:36)
at com.rabbitmq.client.impl.AMQChannel$BlockingRpcContinuation.getReply(AMQChannel.java:443)
at com.rabbitmq.client.impl.AMQChannel.privateRpc(AMQChannel.java:263)
at com.rabbitmq.client.impl.AMQChannel.exnWrappingRpc(AMQChannel.java:136)
... 10 more

Как заставить Flink просто подписаться на очередь RabbitMQ, не пытаясь создать новую?Спасибо всем.

1 Ответ

0 голосов
/ 14 мая 2018

Вы можете написать свой собственный класс, расширяющий RMQSource, и переопределить метод setupQueue, чтобы не создавать очередь

...