内容简介
前言
第1章状态化流处理概述
传统数据处理架构
事务型处理
分析型处理
状态化流处理
事件驱动型应用
数据管道
流式分析
开源流处理的演变
历史回顾
Flink快览
运行首个Flink应用
小结
第2章流处理基础
Dataflow编程概述
Dataflow图
数据并行和任务并行
数据交换策略
并行流处理
延迟和吞吐
数据流上的操作
时间语义
流处理场景下一分钟的含义
处理时间
事件时间
水位线
处理时间与事件时间
状态和一致性模型
任务故障
结果保障
小结
第3章Apache Flink架构
系统架构
搭建Flink所需组件
应用部署
任务执行
高可用性设置
Flink中的数据传输
基于信用值的流量控制
任务链接
事件时间处理
时间戳
水位线
水位线传播和事件时间
时间戳分配和水位线生成
状态管理
算子状态
键值分区状态
状态后端
有状态算子的扩缩容
检查点、保存点及状态恢复
一致性检查点
从一致性检查点中恢复
Flink检查点算法
检查点对性能的影响
保存点
小结
第4章设置Apache Flink开发环境
所需软件
在IDE中运行和调试Flink程序
在IDE中导入书中示例
在IDE中运行Flink程序
在IDE中调试Flink程序
创建Flink Maven项目
小结
第5章DataStream API(1.7版本)
Hello, Flink!
设置执行环境
读取输入流
应用转换
输出结果
执行
转换操作
基本转换
基于KeyedStream的转换
多流转换
分发转换
设置并行度
类型
支持的数据类型
为数据类型创建类型信息
显式提供类型信息
定义键值和引用字段
字段位置
字段表达式
键值选择器
实现函数
函数类
Lambda函数
富函数
导入外部和Flink依赖
小结
第6章基于时间和窗口的算子
配置时间特性
分配时间戳和生成水位线
水位线、延迟及完整性问题
处理函数
时间服务和计时器
向副输出发送数据
CoProcessFunction
窗口算子
定义窗口算子
内置窗口分配器
在窗口上应用函数
自定义窗口算子
基于时间的双流Join
基于间隔的Join
基于窗口的Join
处理迟到数据
丢弃迟到事件
重定向迟到事件
基于迟到事件更新结果
小结
第7章有状态算子和应用
实现有状态函数
在RuntimeContext中声明键值分区状态
通过ListCheckpointed接口实现算子列表状态
使用CheckpointedFunction接口
接收检查点完成通知
为有状态的应用开启故障恢复
确保有状态应用的可维护性
指定算子唯一标识
为使用键值分区状态的算子定义最大并行度
有状态应用的性能及鲁棒性
选择状态后端
选择状态原语
防止状态泄露
更新有状态应用
保持现有状态更新应用
从应用中删除状态
修改算子的状态
可查询式状态
可查询式状态服务的架构及启用方式
对外暴露可查询式状态
从外部系统查询状态
小结
第8章读写外部系统
应用的一致性保障
幂等性写
事务性写
内置连接器
Apache Kafka数据源连接器
Apache Kafka数据汇连接器
文件系统数据源连接器
文件系统数据汇连接器
Apache Cassandra数据汇连接器
实现自定义数据源函数
可重置的数据源函数
数据源函数、时间戳及水位线
实现自定义数据汇函数
幂等性数据汇连接器
事务性数据汇连接器
异步访问外部系统
小结
第9章搭建Flink运行流式应用
部署模式
独立集群
Docker
Apache Hadoop YARN
Kubernetes
高可用性设置
独立集群的HA设置
YARN上的HA设置
Kubernetes的HA设置
集成Hadoop组件
文件系统配置
系统配置
Java和类加载
CPU
内存和网络缓冲
磁盘存储
检查点和状态后端
安全性
小结
第10章Flink和流式应用运维
运行并管理流式应用
保存点
通过命令行客户端管理应用
通过REST API管理应用
在容器中打包并部署应用
控制任务调度
控制任务链接
定义处理槽共享组
调整检查点及恢复
配置检查点
配置状态后端
配置故障恢复
监控Flink集群和应用
Flink Web UI
指标系统
延迟监控
配置日志行为
小结
第11章还有什么?
Flink生态的其他组成部分
用于批处理的DataSet API
用于关系型分析的Table API及SQL
用于复杂事件处理和模式匹配的FlinkCEP
用于图计算的Gelly
欢迎加入社区