首发个人公众号 spark技术分享 , 同步个人网站 coolplayer.net ,未经本人同意,禁止一切转载
Structured Streaming APIS 允许你轻松创建 一种端到端 恰好一次处理保证(exact once) 的 (continuous applications) 持续应用程序, 而且支持高容错性, 用户自己再也不用面对其中涉及的各种麻烦的问题了。
用户可以使用以往熟悉的东西, 比如Spark SQL 中的 DataFrames 和 Datasets 来进行编程, 再也不用面对一些复杂的细节,
上面提到的那么牛逼的功能在一些常见的场景中吸引了越来越多的关注, 比如说 ETL, 复杂的数据格式清洗, 总之可以用在很多地方, Structured Streaming 同时在输入输出方面已经支持了很多常见的组件, 比如 Kafka, HDFS, S3, RDBMS 等,
在这篇博客中, 我会讲到Structured Streaming 怎么 和 Kafka 进行端到端的整合, 从简单处理到复杂的窗口 ETL, 同时会根据自己的需求, 把数据输出到 内存, console, 文件, 数据库, 或者 kafka 中, 在输出到文件的 例子中, 我会把新的数据 覆盖已经存在的分区表中。

连接kafka
假设你现在已经有一个正常服务的 Kafka 集群了, 你现在需要做的就是使用 Structured Streaming 从一个 topic 中消费处理消息, 现在已经有现成的 Structured Streaming 到 Kafka 的连接器了, 所以很简单地就可以对消息进行流式读取

这里你可以设置很多参数来控制这个读取过程, 你可以在 这里 参考一下这些参数,
我们快速看下 我们启动的 streamingInputDF DataFrame 的 schema,

包括 key, value, topic, partition, offset, timestamp and timestampType 等字段, 我可以选择抽取我们需要的字段, value
这个字段包含我们真正的消息数据, 其他的都是消息附带的元数据, timestamp 这个字段代表消息到达的时间戳, 这里注意 不要和 消息里面事件产生的时间的时间戳 搞混了。 一个是 到达 kafka的时间, 一个是事件产生的时间, 当然 事件产生的时间 才是我们真正要关心的。
流式 ETL
数据的读取已经开始, 我们现在要对数据进行数据的抽取和聚合, 注意, 到这里 streamingInputDF 是一个 DataFrame类型, 因为 DataFrames 是一个 没有类型的数据集, 我们需要做一些处理,
假如我们的数据是 标准的ISP 的日志, 我们的例子就如下所示:

现在就可以进行一些有趣的分析了, 比如每个 zipcode 有多少用户, 以及 这些用户从 什么 ISP来的, 我们就可以创建 一个运营的指标面板, 给其他需要的部门了

我们刚才对json 数据进行解析, 然后分组聚合, 这些操作都是准实时的,而且随着 kafka 中数据的更新,我们的结果也是持续更新的。 每当有新的数据到来, 就会触发一次增量的更新,
窗口操作
我们上面已经实现了, 对json数据的持续的解析聚合, 现在我们想实现一个这样的需求: 从小时开始的 2分钟, 进行 10分钟窗口的流量进行聚合, 每5分钟窗口滑动一次,
我们的 json 数据中包含 hittime 作为事件发生的时间戳,
下图中的 饼图中代表着 十分钟的窗口聚合的指标

输出结果
目前为止, 我们能看到自动更新的指标, 如果我们想进一步控制输出, 输出的时候可以用很多参数来控制, 比如我们想对应用 debug, 那我们肯定希望数据打印在控制台上, 如果我们想同时对数据进行一些交互式的查询, 那我们放在内存中更好点, 当然你也可以把数据输出到数据库,文件,甚至kafka中。
内存
在这个场景中, 我们把数据放在内存中的一张表中, 用户就可以对这张表进行 sql 查询, 表的名字也可以指定, 我们还是使用 streamingSelectDF 的例子,

现在你就可以对你感兴趣的指标进行分析了, 这个时候你是在一个动态无限大的表上进行查询的
控制台
输出打印到了控制台上

文件
这种方式适用于需要持久化的场景, 不像输出到内存和 控制台中, 这种方式是需要容错的, 所以要设置一个检查点目录, 用于容错需要的 状态存储,


一旦数据保存后, Spark 就可以向查询其他数据集一样, 查询这些文件,

还有一个优势, 就是我们可以使用任意字段来进行动态分区, 上面这个例子中, 我们可以使用 ‘zipcode’ and ‘day’ 字段来进行分区, 这样就可以 skip 大量不需要关心的数据, 比如我们需要聚合今天的, 那么只需要消费今天的分区就好了

现在我们用‘zip’ by ‘day’ 来对流中的数据进行分区,

我们看下输出目录中的文件,

这些分区后的数据可以直接 在 datasets and DataFrames中使用, 如果设置 一个表指向这些文件, Spark SQL 就可以直接查询,

使用这种方法的一个注意事项是,必须将分区添加到表中才能访问其下的数据集。

分区可以事先设置, 一旦有文件, 就立即可用了。

现在你就可以在这些自动更新的并且持久化了的数据上进行一些分析了
数据库
经常我们需要把 结果输出到外部数据中, 比如 mysql, 现在, Structured Streaming API 还不支持外部数据库,如果支持,当然就可以简化为 .format(“jdbc”).start(“jdbc:mysql/..”) 。
现在我们需要 用 foreach sink 来完成这个需求, 我们创建一个 自定义的 JDBC Sink 继承 ForeachWriter 并实现方法,

现在我们使用 JDBCSink:

当 一批数据完成了, 分组聚合的 counts, 就可以 插入/覆盖 mysql 中的数据了。

Kafka
就像写数据库一样, 现在 Structured Streaming API 还不支持直接写入 kafka, 那么我们还是 使用 KafkaSink 继承 _ForeachWriter 来实现

现在我们可以使用这个 writer:

你可以看到我们把消息打入了 topic2, 这个demo 中, 我们打入的是 每一批数据 数据到达后, 增量更新后的 zipcode:count, 还要注意的是, streaming 的dashboard 中, 我们可以看到,现在到达速率, 和处理速率, 以及批处理时间, 这让 我们监控系统的性能很方便。
我们使用 kafka 消费看下结果

在这个demo 中,我们使用的 update
输出模式, 一批消息被消费后, 只有更新了的 zipcodes才会被打到 kafka, 没有变化的 Zipcodes 则不会, 当然你也可以 使用 complete 模式, 每次所有的 zipcode 都会被打到kafka。
小结
这里, 我简单的 测试了 Structured Streaming 和 Kafka 的整合, 这里展示了一些不同的输入和输出方式, 其实 对于其他不同的数据源, 流程都差不多, 像 sockets, directory等, 如果你的数据源 是 socket, 只要简单的换一下就好了。
原创精品,首发个人公众号 spark技术分享 , 同步个人网站 coolplayer.net ,未经本人同意,禁止一切转载
欢迎关注公众号 spark技术分享





