T112019-数据智能技术峰会-Flink在数据分析中的应用-2019.11.25-24页_3mb
报告摘要
Flink 在 TalkingData SaaS 分析中的应用总结
核心内容
Flink 在 TalkingData SaaS 分析系统中被广泛采用,以解决传统流处理框架在性能、扩展性和资源管理上的痛点。随着业务量的增长,TalkingData 逐步从 Jetty 服务和自研 etl-framework 过渡到 Flink,并在多个阶段进行了优化与演进。
主要观点
- 流处理背景与痛点:TalkingData 早期使用 Jetty 服务和自研 etl-framework,但存在扩展性差、性能瓶颈以及资源分配不均的问题。
- Flink 优势:Flink 提供了更强大的流处理能力,支持 exactly-once 和 at-least-once 语义,具备良好的内存管理和 SQL 支持,适用于更复杂的业务场景。
- 演进路线:TalkingData 在 Flink 的使用上经历了从 standalone 集群到 Flink on Yarn 的演进,以提高资源利用率和作业隔离性。
- 性能优化:通过优化 Job 部署结构、减少网络传输、合理配置 TaskManager 和 Container 粒度等方式,提升了系统的吞吐量和稳定性。
- 序列化优化:通过使用 POJOs 和 TypeInformation,减少了序列化对 CPU 的影响,提升了处理效率。
关键信息
1. TalkingData 流处理的背景与痛点
- Jetty 服务:难以扩展和维护,性能不足。
- 自研 etl-framework:无法完整表达 DAG,容错机制不足,性能仍有问题。
- 新流处理系统:旨在满足更大的业务量和更复杂的业务场景。
2. Flink 在 TalkingData SaaS 分析中的演进路线
- Standalone Cluster:
- 2017.4~2017.6:单集群,48 cores,日均处理 42 亿数据包。
- 2017.7~2017.8:单集群,120 cores,日均处理 46 亿数据包。
- 2017.9~2017.12:单集群,264 cores,日均处理 46 亿数据包。
- 2018.1~2018.12:拆分为 2 个集群,456 cores,日均处理 46 亿数据包。
- Flink on Yarn:
- 2019.7:日均处理 63 亿数据包,峰值 9.5 万 package/s,80 万 events/s。
- 支持多租户分发、调度、监控和流式与批处理队列的混合部署。
3. 实践经验
3.1 Job 阻塞与网络栈优化
- 问题:Job 阻塞导致吞吐量下降。
- 解决方案:
- 将 operator chain 在一起,减少网络传输和序列化反序列化的开销。
- 使用 Flink 1.5 及以上版本。
- 优化 Buffer 管理,避免 BufferPool 中 Buffer 不足导致的阻塞。
3.2 资源的 balance 与 isolation
- 问题:Standalone 集群中资源分配不均,导致 TM 之间资源浪费。
- 解决方案:
- 将 TaskManager 粒度变小,每台机器部署多个实例。
- 将大业务 job 隔离到不同的集群。
- 采用 Flink on Yarn 模式,实现更细粒度的资源管理和作业隔离。
- Container 拆解粒度不宜过小,建议 2 core 4G 为宜。
3.3 序列化与反序列化优化
- 问题:序列化成为 CPU 抽样的热点,影响性能。
- 解决方案:
- 使用 POJOs(Java Bean)和 TypeInformation 提升序列化效率。
- 显式调用 returns 方法,触发 Flink 的 Type Hint。
- 注册数据类及其子类、字段类型,使 Flink 更高效地处理类型信息。
- 自定义序列化器以进一步优化性能。
总结与展望
- 当前成果:Flink 已稳定支持 TalkingData 分析线,日均处理 63 亿数据包,峰值 9.5 万 package/s 和 80 万 events/s。
- 未来展望:可进一步探索将更复杂的业务场景迁移到 Flink 上,甚至支持 Batch Job,以实现更全面的数据处理能力。
展开完整摘要
试读结束,高清完整版pdf/doc/ppt,请点下载