【发布时间】:2017-01-20 10:08:34
【问题描述】:
首先,这与Kafka consuming the latest message again when I rerun the Flink consumer 非常相似,但又不一样。该问题的答案似乎无法解决我的问题。如果我错过了该答案中的某些内容,请改写答案,因为我显然错过了某些内容。
但问题完全相同——Flink(kafka 连接器)重新运行它在关闭之前看到的最后 3-9 条消息。
我的版本
Flink 1.1.2
Kafka 0.9.0.1
Scala 2.11.7
Java 1.8.0_91
我的代码
import java.util.Properties
import org.apache.flink.streaming.api.windowing.time.Time
import org.apache.flink.streaming.api.scala._
import org.apache.flink.streaming.api.CheckpointingMode
import org.apache.flink.streaming.connectors.kafka._
import org.apache.flink.streaming.util.serialization._
import org.apache.flink.runtime.state.filesystem._
object Runner {
def main(args: Array[String]): Unit = {
val env = StreamExecutionEnvironment.getExecutionEnvironment
env.enableCheckpointing(500)
env.setStateBackend(new FsStateBackend("file:///tmp/checkpoints"))
env.getCheckpointConfig.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE)
val properties = new Properties()
properties.setProperty("bootstrap.servers", "localhost:9092");
properties.setProperty("group.id", "testing");
val kafkaConsumer = new FlinkKafkaConsumer09[String]("testing-in", new SimpleStringSchema(), properties)
val kafkaProducer = new FlinkKafkaProducer09[String]("localhost:9092", "testing-out", new SimpleStringSchema())
env.addSource(kafkaConsumer)
.addSink(kafkaProducer)
env.execute()
}
}
我的 SBT 依赖项
libraryDependencies ++= Seq(
"org.apache.flink" %% "flink-scala" % "1.1.2",
"org.apache.flink" %% "flink-streaming-scala" % "1.1.2",
"org.apache.flink" %% "flink-clients" % "1.1.2",
"org.apache.flink" %% "flink-connector-kafka-0.9" % "1.1.2",
"org.apache.flink" %% "flink-connector-filesystem" % "1.1.2"
)
我的过程
(3 个终端)
TERM-1 start sbt, run program
TERM-2 create kafka topics testing-in and testing-out
TERM-2 run kafka-console-producer on testing-in topic
TERM-3 run kafka-console-consumer on testing-out topic
TERM-2 send data to kafka producer.
Wait for a couple seconds (buffers need to flush)
TERM-3 watch data appear in testing-out topic
Wait for at least 500 milliseconds for checkpointing to happen
TERM-1 stop sbt
TERM-1 run sbt
TERM-3 watch last few lines of data appear in testing-out topic
我的期望
当系统中没有错误时,我希望能够打开和关闭 flink,而无需重新处理在先前运行中成功完成流的消息。
我的修复尝试
我已将调用添加到setStateBackend,认为可能是默认内存后端没有正确记住。这似乎没有帮助。
我已经删除了对enableCheckpointing 的调用,希望在 Flink 和 Zookeeper 中可能有一个单独的机制来跟踪状态。这似乎没有帮助。
我使用了不同的接收器,RollingFileSink,print();希望这个错误可能在kafka中。这似乎没有帮助。
我已经回滚到 flink(和所有连接器)v1.1.0 和 v1.1.1,希望这个 bug 可能在最新版本中。这似乎没有帮助。
我已将zookeeper.connect 配置添加到属性对象中,希望关于它仅在 0.8 中有用的评论是错误的。这似乎没有帮助。
我已将检查点模式明确设置为 EXACTLY_ONCE(好主意 drfloob)。这似乎没有帮助。
我的请求
救命!
【问题讨论】:
-
只是为了好玩,尝试明确设置 EXACTLY_ONCE。 env.getCheckpointConfig.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE)
-
我遇到了完全相同的问题,在启动 flink 流作业后再次消费相同的事件。也许保存检查点时偏移量没有正确增加?
-
显式检查点模式不起作用。我已经更新了帖子。不过是个好主意。
标签: duplicates apache-kafka apache-flink flink-streaming