Я пытаюсь написать доказательство концепции, которое берет сообщения от Kafka, преобразует их с помощью Beam на Flink, а затем отправляет результаты в другую тему Kafka.
Я использовал KafkaWindowedWordCountExample в качестве отправной точки, и он делает первую часть того, что я хочу сделать, но он выводит в текстовые файлы, а не в Kafka. FlinkKafkaProducer08 выглядит многообещающе, но я не могу понять, как подключить его к конвейеру. Я думал, что он будет обернут UnboundedFlinkSink или чем-то подобным, но, похоже, этого не существует.
Любые советы или мысли о том, что я пытаюсь сделать?
Я запускаю последний луч-инкубатор (по состоянию на прошлый вечер с Github), Flink 1.0.0 в кластерном режиме и Kafka 0.9.0.1, все на Google Compute Engine (Debian Jessie).