【问题标题】:Flink task managers are not processing data after restartFlink 任务管理器重启后不处理数据
【发布时间】:2020-12-04 17:31:46
【问题描述】:

我是 flink 新手,我部署了我的 flink 应用程序,它基本上执行简单的模式匹配。它部署在 1 个 JM 和 6 个 TM 的 Kubernetes 集群中。我每 10 分钟向 eventthub 主题发送大小为 4.4k 和 200k 的消息并执行负载测试。我添加了重启策略并检查指向如下,我没有在我的代码中明确使用任何状态,因为不需要它

 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

 // start a checkpoint every 1000 ms
 env.enableCheckpointing(interval, CheckpointingMode.EXACTLY_ONCE);
 // advanced options:
 // make sure 500 ms of progress happen between checkpoints
 env.getCheckpointConfig().setMinPauseBetweenCheckpoints(1000);
 // checkpoints have to complete within one minute, or are discarded
 env.getCheckpointConfig().setCheckpointTimeout(120000);
 // allow only one checkpoint to be in progress at the same time
 env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);
 // enable externalized checkpoints which are retained after job cancellation
 env.getCheckpointConfig().enableExternalizedCheckpoints(CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
 // allow job recovery fallback to checkpoint when there is a more recent savepoint
 env.getCheckpointConfig().setPreferCheckpointForRecovery(true);

 env.setRestartStrategy(RestartStrategies.fixedDelayRestart(
         5, // number of restart attempts
         Time.of(5, TimeUnit.MINUTES) // delay
 ));

最初我遇到网络缓冲区的 Netty 服务器问题,我点击此链接 https://ci.apache.org/projects/flink/flink-docs-release-1.11/ops/config.html#taskmanager-network-memory-floating-buffers-per-gate flink 网络和堆内存优化并应用以下设置,一切正常

taskmanager.network.memory.min: 256mb
taskmanager.network.memory.max: 1024mb
taskmanager.network.memory.buffers-per-channel: 8
taskmanager.memory.segment-size: 2mb
taskmanager.network.memory.floating-buffers-per-gate: 16
cluster.evenly-spread-out-slots: true
taskmanager.heap.size: 1024m
taskmanager.memory.framework.heap.size: 64mb
taskmanager.memory.managed.fraction: 0.7
taskmanager.memory.framework.off-heap.size: 64mb
taskmanager.memory.network.fraction: 0.4
taskmanager.memory.jvm-overhead.min: 256mb
taskmanager.memory.jvm-overhead.max: 1gb
taskmanager.memory.jvm-overhead.fraction: 0.4

但我有以下两个问题

  1. 如果任何任务管理器由于任何故障而重新启动,则任务管理器将成功重新启动并注册到作业管理器,但在重新启动的任务管理器不执行任何数据处理后,它将处于空闲状态。这是正常的 flink 行为还是我需要添加任何设置以使任务管理器重新开始处理。

  2. 抱歉,如果我的理解有误,请纠正我,flink 在我的代码中有重启策略,我限制了 5 次重启尝试。如果我的 flink 作业没有成功克服任务失败会发生什么整个 flink 作业将保持在空闲状态,我必须手动重新启动作业,或者我可以添加任何机制来重新启动我的作业,即使它超过了重新启动的限制工作尝试。

  3. 是否有任何文档可以根据我的系统接收数据的数据大小和速率来计算我应该分配给 flink 作业集群的内核和内存数量?

  4. 有没有关于 flink CEP 优化技术的文档?

  5. 这是我在作业管理器中看到的错误堆栈跟踪

  1. 在模式匹配之前,我在作业管理器日志中看到以下错误

    原因:org.apache.flink.runtime.io.network.netty.exception.RemoteTransportException:远程任务管理器“/10.244.9.163:46377”意外关闭连接。这可能表明远程任务管理器丢失了。 在 org.apache.flink.runtime.io.network.netty.CreditBasedPartitionRequestClientHandler.channelInactive(CreditBasedPartitionRequestClientHandler.java:136) 在 org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:257) 在 org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:243) 在 org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.fireChannelInactive(AbstractChannelHandlerContext.java:236) 在 org.apache.flink.shaded.netty4.io.netty.handler.codec.ByteToMessageDecoder.channelInputClosed(ByteToMessageDecoder.java:393) 在 org.apache.flink.shaded.netty4.io.netty.handler.codec.ByteToMessageDecoder.channelInactive(ByteToMessageDecoder.java:358) 在 org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:257) 在 org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:243) 在 org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.fireChannelInactive(AbstractChannelHandlerContext.java:236) 在 org.apache.flink.shaded.netty4.io.netty.channel.DefaultChannelPipeline$HeadContext.channelInactive(DefaultChannelPipeline.java:1416) 在 org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:257) 在 org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:243) 在 org.apache.flink.shaded.netty4.io.netty.channel.DefaultChannelPipeline.fireChannelInactive(DefaultChannelPipeline.java:912) 在 org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannel$AbstractUnsafe$8.run(AbstractChannel.java:816) 在 org.apache.flink.shaded.netty4.io.netty.util.concurrent.AbstractEventExecutor.safeExecute(AbstractEventExecutor.java:163) 在 org.apache.flink.shaded.netty4.io.netty.util.concurrent.SingleThreadEventExecutor.runAllTask​​s(SingleThreadEventExecutor.java:416) 在 org.apache.flink.shaded.netty4.io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:515) 在 org.apache.flink.shaded.netty4.io.netty.util.concurrent.SingleThreadEventExecutor$5.run(SingleThreadEventExecutor.java:918) 在 org.apache.flink.shaded.netty4.io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74) 在 java.lang.Thread.run(Thread.java:748)

提前谢谢,请帮我解决我的疑惑

【问题讨论】:

    标签: apache-flink flink-streaming flink-cep


    【解决方案1】:

    各点:

    如果您的模式涉及匹配时间序列(例如,“A 后跟 B”),那么您需要状态来执行此操作。 Flink 的大部分 source 和 sinks 在内部也使用 state 来记录 offset 等,如果你关心完全一次性的保证,这个 state 需要被检查点。如果模式是动态流式传输的,那么您也需要将模式存储在 Flink 状态中。

    代码中的某些 cmets 与配置参数不匹配:例如,“500 ms 的进度”与 1000、“检查点必须在一分钟内完成”与 120000。另外,请记住该部分您从中复制这些设置的文档不是推荐最佳实践,而是说明如何进行更改。特别是,env.getCheckpointConfig().setPreferCheckpointForRecovery(true); 是个坏主意,而且该配置选项可能不存在。

    您在 config.yaml 中的一些条目与此相关。 taskmanager.memory.managed.fraction 相当大(0.7)——这仅在您使用 RocksDB 时才有意义,因为托管内存没有其他用于流式传输的目的。而taskmanager.memory.network.fraction和taskmanager.memory.jvm-overhead.fraction都很大,这三个分数之和是1.5,没有意义。

    一般来说,默认网络配置适用于各种部署场景,并且需要调整这些设置的情况很少见,但大型集群除外(此处不是这种情况)。你遇到过什么样的问题?

    至于你的问题:

    1. TM 发生故障并恢复后,TM 应自动从最近的检查点恢复处理。要诊断为什么没有发生这种情况,我们需要更多信息。要获得正确处理此问题的部署经验,您可以尝试Flink Operations Playground。

    2. 一旦配置的重启策略执行完毕,作业将失败,Flink 将不再尝试恢复该作业。当然,如果你想要更复杂的东西,你可以在 Flink 的 REST API 之上构建自己的自动化。

    3. 关于容量规划的文档?不,不是真的。这通常是通过反复试验来解决的。不同的应用程序往往以难以预料的方式有不同的要求。诸如您选择的序列化程序、状态后端、keyBys 的数量、源和接收器、密钥偏差、水印等都会产生重大影响。

    4. 关于优化 CEP 的文档?不,对不起。要点是

      • 尽你所能限制匹配;避免必须无限期保持状态的模式
      • getEventsForPattern 可能很贵

    【讨论】:

    • 感谢您的回复,我在模式匹配中所做的是,我得到一个 json 流,我需要在其中验证一组规则并广播所有规则,我正在使用广播处理函数我的 json 流是非键控流。如果 json 满足规则,那么我将作为有效的 json 发送到下游。我在 flink 文档中阅读了很多关于状态的信息,并且可以在示例中看到它们将总和值存储在状态描述符中。在我的情况下,我没有对我的有效负载执行任何计算,您能否建议我如何在这种情况下实现状态。
    • 我没有使用 CEP 库来执行模式匹配。
    • 我假设您正在使用 CEP,或者需要状态来匹配涉及有序事件序列的模式。如果您只是匹配单个事件,那么您需要的唯一状态可能是源和接收器的状态,以及用于存储模式的广播状态。
    • 我重写了第一段以尝试更好地适应您的用例。
    • 很抱歉再次打扰您,我有最后一个问题,我正在接收来自多个来源的信号,并且我在广播状态下维护所有规则。当流来自特定源时,我将在广播状态下过滤该源的规则,例如 Map ruleList = ctx.getBroadcastState(broadcastDescriptor).get(source) 并且我正在遍历 ruleList 并应用规则我的流,对于一个特定的来源,我有 60 条规则,因为模式匹配模块正在关闭。有没有办法避免迭代并执行模式匹配?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多