Flink在数据分析中的应用_24页_3mb
报告摘要
Flink 在 TalkingData 数据分析中的应用总结
核心内容
Flink 在 TalkingData 的 SaaS 分析系统中被广泛应用,主要用于处理海量数据流。TalkingData 在流处理服务的演进过程中经历了多个阶段,从早期的 Jetty 服务到自研的 etl-framework,再到基于 Flink 的流处理系统。随着业务规模的扩大,原有的系统在扩展性、性能和资源管理方面暴露出诸多问题,促使 TalkingData 采用 Flink 作为核心流处理引擎。
主要观点
- 流处理服务演进:TalkingData 流处理服务从 Jetty 到自研 etl-framework,再到 Flink,逐步解决了扩展性差、性能瓶颈、资源分配不均等问题。
- Flink 的优势:Flink 在性能、语义支持(exactly-once)、SQL 支持、资源管理等方面优于其他流处理框架(如 Storm),并支持批处理作为流处理的特例。
- 资源隔离与调度:通过将任务隔离到不同的集群,TalkingData 有效解决了资源竞争和性能瓶颈问题,提升了系统的稳定性和吞吐量。
- 序列化优化:Flink 提供了丰富的序列化机制,包括 TypeInformation、POJOs 和自定义序列化器,以减少 CPU 使用和内存开销。
关键信息
1. TalkingData 流处理服务演进
- Jetty 服务:存在不易扩展与维护、性能问题。
- 自研 etl-framework:无法完整表达 DAG、容错机制不足、性能问题。
- 新的流处理系统:满足更大的业务量和更复杂的业务场景。
2. Flink 在 TalkingData SaaS 分析中的演进路线
- Standalone Cluster:初期部署单集群,随着数据量增长,出现资源分配不均和 Job 阻塞问题,最终拆分为多个集群。
- Flink on Yarn:通过 Yarn 集群实现多租户分发与调度,提升资源利用率和系统稳定性。
- 资源分配:通过调整 TaskManager 的粒度和隔离大业务 Job 到不同集群,优化资源平衡与隔离。
3. 实践经验
3.1 Job 阻塞与网络栈优化
- 问题:Job 阻塞导致吞吐量急剧下降。
- 解决方案:
- 尽可能将 Operator 链在一起,减少网络传输和序列化反序列化的开销。
- 使用 Flink 1.5 及之后的版本,优化网络栈性能。
3.2 资源的 Balance 与 Isolation
- 问题:Standalone 集群中 TaskManager 的资源分配不均,导致部分节点空闲,影响整体性能。
- 解决方案:
- 将 TaskManager 粒度变小,部署多个实例,每个实例持有较少的 slot。
- 将大业务 Job 隔离到不同的集群,减少资源竞争。
- 采用 Flink on Yarn 实现更细粒度的资源管理和调度。
3.3 序列化与反序列化优化
- 问题:序列化成为 CPU 抽样的热点,对性能损耗较大。
- 解决方案:
- 使用 TypeInformation 提供更详细的类型信息,优化序列化方式。
- 推荐使用 POJOs(Java Bean)来提升性能。
- 显式调用
returns方法触发 Flink 的 Type Hint。 - 通过
registerType方法注册自定义类型及其子类,提升序列化效率。 - 自定义序列化器以进一步优化性能。
4. 总结与展望
- 当前状态:Flink 已经稳定支持 TalkingData 分析线,日均处理 63 亿个 package,峰值可达 9.5 万 package/s 和 80 万 events/s。
- 未来展望:可以进一步探索将更复杂的业务迁移到 Flink 上,甚至支持批处理 Job。
关键数据对比
| 阶段 | 日均数据量 | 峰值 package/s | 峰值 events/s | 集群规模(core) |
|---|---|---|---|---|
| Standalone Cluster | 42 亿 | 2.8 万 | 40 万 | 120 核 |
| Flink on Yarn | 15 亿 | 9.5 万 | 80 万 | 432 核 |
总结
Flink 在 TalkingData 的数据分析系统中发挥了重要作用,解决了早期流处理系统的诸多痛点。通过不断优化资源管理和网络栈性能,TalkingData 实现了更高的吞吐量和更稳定的运行环境。同时,Flink 的序列化机制也得到了有效利用,提升了整体处理效率。未来,Flink 有望在 TalkingData 的更多业务场景中发挥作用,进一步推动数据分析能力的提升。
展开完整摘要
试读结束,高清完整版pdf/doc/ppt,请点下载