改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

Flink实时数仓保障体系构建与优化策略

ckckck2025年1月10日13 浏览

时效性保证

  1. 查看Kafka延迟监控:通过Flink消费上游的lag(如消费Kafka的lag情况)来监控延迟。
  2. 分层和时延的平衡:在复用性和链路过长之间做好平衡和取舍。
  3. 乱序数据处理:处理乱序数据以确保数据顺序正确。
  4. 提前压测:应对流量高峰期,特别是大促场景下,提前做好资源保障和任务优化。
  5. 设置延时基线:通过优化程序代码、资源、解决倾斜与反压等问题,将延迟控制在基线之内。
  6. 指标监控:监控任务failover情况、checkpoint指标、GC情况、作业反压等,出现异常时及时告警。
  7. Flink链路延迟监控:通过LatencyMarker观测延迟情况。

质量考量

数据一致性

  1. 正确性实时计算端到端的一致性:通过输出幂等方式保障数据正确性,要求存储介质支持重写。对于不支持幂等的存储(如DWD层的Kafka),可以使用row_number()语法去重。
  2. 离线与实时的一致性:确保使用相同的数据源和加工业务逻辑。

数据完整性

确保数据从源头到前端展示的完整性,避免因加工逻辑、存储异常或前端展现异常导致数据丢失。例如:

  1. 数据源层背压:可能导致数据源头(MQ, Kafka)消息积压,严重时导致数据丢失。
  2. 数据处理层加工错误:未按需求加工数据,导致目标有效数据丢失。
  3. 数据存储层容量写满:新数据无法写入,导致数据丢失。
  4. 数据加工正确性、及时性、快速恢复性:构成数据完整性的关键要素。

数据加工正确性

确保源数据按业务需求加工成目标有效数据,并根据不同维度计算展示指标。例如:

  1. 数据源层原始数据过滤:过滤掉不需要的联盟点击数据,补齐目标联盟的点击数据。
  2. 业务层计算:根据媒体、账号、计划、单元等维度计算点击总量。

数据加工及时性

控制数据从产生到前端展示的时间在合理范围内。

数据快速恢复性

确保数据在流转路径中因异常中断后能快速恢复,且恢复过程正确无误。例如:

  1. 数据处理层性能问题:解决后数据积压问题逐步缓解。
  2. 消费程序崩溃:重启后数据消费恢复正常。

数据可监控性

确保数据流转路径中关键节点的状态可有效监控。

数据高可用性

确保数据不会因灾难性问题丢失,实时数据消费应用集群和存储集群需具备主备和容灾能力。

稳定考量

  1. 任务压测:提前压测应对流量高峰期,做好资源保障和任务优化。
  2. 任务分级:根据任务影响面大小和数据使用方划分保障等级,高优先级任务需优先响应。
  3. 指标监控:监控任务failover情况、checkpoint指标、GC情况、作业反压等,出现异常时告警。
  4. 高可用HA:确保实时Pipeline链路高可用,支持数据备份和重演机制。
  5. SLA保障:支持动态扩容和数据处理流程自动漂移。
  6. 弹性反脆弱:基于规则和算法的资源弹性伸缩,支持事件触发动作引擎的失效处理。
  7. 监控预警:集群设施、物理管道、数据逻辑层面的多方面监控预警能力。
  8. 自动运维:捕捉并存档缺失数据和处理异常,具备定期自动重试机制修复问题数据。
  9. 上游元数据变更抗性:上游业务库需兼容元数据变更,实时Pipeline处理显式字段。

成本考量

  1. 人力成本:通过支持数据应用平民化降低人才人力成本。
  2. 资源成本:通过支持动态资源利用降低静态资源占用造成的资源浪费。
  3. 运维成本:通过支持自动运维、高可用、弹性反脆弱等机制降低运维成本。
  4. 试错成本:通过支持敏捷开发、快速迭代降低试错成本。

敏捷考量

敏捷大数据是一整套理论体系和方法学,意味着配置化、SQL化、平民化。

管理考量

数据管理重点关注元数据管理和数据安全管理。在现代数仓多数据存储选型的环境下,统一管理元数据和数据安全是一个挑战,需在实时Pipeline各个环节平台内置支持,并支持对接外部统一的元数据管理平台和数据安全策略。

本文主要讲述了在实时数仓建设中常见的任务保障方法和理论,具体实践可能有所不同,但总体思路值得借鉴。

end