行业资讯
Flink实战(10)-checkpoint容错保证
Flink实战(10)-checkpoint容错保证0 前言程序在 Flink 集群运行某个算子因为某些原因出现故障如何处理在故障恢复后如何保证数据状态和故障发生之前的数据状态一致?1 什么是 checkpoint(检查点)?Checkpoint 能生成快照(Snapshot)。若 Flink 程序崩溃重新运行程序时可以有选择地从这些快照进行恢复。Checkpoint 是 Flink 可靠性的基石。2 Checkpoint V.S StateState 指某个算子的数据状态保存在堆内存Checkpoint 指所有算子的数据状态持久化保存3 什么是savepoint(保存点)?基于 checkpoint 机制的快照。4 Checkpoint V.S SavepointCheckpoint 是 自动容错恢复机制Savepoint 某个时间点的全局状态镜像Checkpoint 是 Flink 系统行为 。Savepoint 是用户触发Checkpoint 默认程序删除。Savepoint 会一直保存5 数据流快照最简单的流程暂停处理新流入数据将新数据缓存起来将算子任务的本地状态数据拷贝到一个远程的持久化存储上继续处理新流入的数据包括刚才缓存起来的数据6 Flink slot 和并行度设置合理的并行度能够加快数据的处理Flink 每个算子都可以设置并行度Slot 使得 taskmanager 具有并发执行的能力Flink 任务和子任务从 Source 到 sink每当并行度发生变化或者数据分组( keyBy)就会产生任务。一个任务的并行度为 N就会有 N 个子任务。7 Checkpoint 分布式快照流程第1步要实现分布式快照最关键的是能够将数据流切分。Flink 中使用 Checkpoint Barrier(检查点分割线)来切分数据流当 Source 子任务收到 Checkpoint 请求该算子会对自己的数据状态保存快照。向自己的下一个算子发送 Checkpoint Barrier下一个算子只有收到上一个算子广播过来的 Checkpoint Barrier才进行快照保存。第2步当 Sink 算子已经收到所有上游的 Checkpoint Barrie 时进行以下 2 步操作保存自己的数据状态并直接通知检查点协调器检查点协调器在收集所有的 task 通知后就认为这次的 Checkpoint 全局完成了。下游算子有多个数据流输入啥时才 checkpoint这就涉及到Barrie对齐机制保证了 Checkpoint 数据状态的精确一致。第1步下一个算子某个通道接收了第一个ID为n的 Checkpoint Barrie这个算子其他通道的ID 为n的 Checkpoint Barrie 还没到达第2步该算子将第一个ID为n的 Checkpoint Barrie 缓存该个算子继续处理其他通道的ID为n的 Checkpoint Barrie第3步:该个算子所有通道的ID 为n的 Checkpoint Barrie 到达后该算子执行快照不进行 Barrier 对齐可以吗8 Checkpoint咋保证数据状态的一致性Flink内置的数据状态一致性端到端的数据状态一致性Flink 系统内部的数据状态一致性AT-MOST-ONCE(最多一次已废除)发生故障可能会丢失数据AT-LEAST-ONCE(至少一次)发生故障可能会有重复数据。EXACTLY-ONCE(精确一次)发生故障能保证不丢失数据也没有重复数据图片KafkaSink总共支持三种不同的语义保证DeliveryGuarantee。对于DeliveryGuarantee.AT_LEAST_ONCE和DeliveryGuarantee.EXACTLY_ONCEFlink checkpoint 必须启用。默认情况下KafkaSink使用DeliveryGuarantee.NONE。DeliveryGuarantee.NONE不提供任何保证消息有可能会因 Kafka broker 的原因发生丢失或因 Flink 的故障发生重复。DeliveryGuarantee.AT_LEAST_ONCE: sink 在 checkpoint 时会等待 Kafka 缓冲区中的数据全部被 Kafka producer 确认。消息不会因 Kafka broker 端发生的事件而丢失但可能会在 Flink 重启时重复因为 Flink 会重新处理旧数据。DeliveryGuarantee.EXACTLY_ONCE: 该模式下Kafka sink 会将所有数据通过在 checkpoint 时提交的事务写入。因此如果 consumer 只读取已提交的数据参见 Kafka consumer 配置isolation.level在 Flink 发生重启时不会发生数据重复。然而这会使数据在 checkpoint 完成时才会可见因此按需调整 checkpoint 间隔。请确认事务 ID 的前缀transactionIdPrefix对不同的应用是唯一的以保证不同作业的事务 不会互相影响此外强烈建议将 Kafka 的事务超时时间调整至远大于 checkpoint 最大间隔 最大重启时间否则 Kafka 对未提交事务的过期处理会导致数据丢失。9 Data Source 和 Sink 的容错保证当程序出现错误的时候Flink 的容错机制能恢复并继续运行程序。这种错误包括机器硬件故障、网络故障、瞬态程序故障等。只有当 source 参与快照机制Flink 才能保证对自定义状态的精确一次更新。下表列举了 Flink 与其自带连接器的状态更新的保证。SourceGuaranteesNotesApache Kafka精确一次根据你的版本用恰当的 Kafka 连接器Amazon Kinesis Data Streams精确一次RabbitMQ至多一次 (v 0.10) / 精确一次 (v 1.0)Google PubSub至少一次Collections精确一次Files精确一次Sockets至多一次为保证端到端精确一次的数据交付在精确一次的状态语义上更进一步sink需要参与 checkpointing 机制。下表列举了 Flink 与其自带 sink 的交付保证假设精确一次状态更新。SinkGuaranteesNotesElasticsearch至少一次Opensearch至少一次Kafka producer至少一次 / 精确一次当使用事务生产者时保证精确一次 (v 0.11)Cassandra sink至少一次 / 精确一次只有当更新是幂等时保证精确一次Amazon DynamoDB至少一次Amazon Kinesis Data Streams至少一次Amazon Kinesis Data Firehose至少一次File sinks精确一次Socket sinks至少一次Standard output至少一次Redis sink至少一次
郑州网站建设
网页设计
企业官网