如何基于 Apache Pulsar 和 Spark 进行批流一体的弹性资料处理?
在大规模并行资料分析领域,AMPLab 的‘One stack to rule them all’提出用 Apache Spark 作为统一的引擎支援批处理、流处理、互动查询和机器学习等常见的资料处理场景。
批流现状
在大规模并行资料分析领域,AMPLab 的‘One stack to rule them all’提出用 Apache Spark 作为统一的引擎支援批处理、流处理、互动查询和机器学习等常见的资料处理场景。 2017 年 7 月,Spark 2.2.0 版本正式推出的 Spark structured streaming 将 Spark SQL 作为流处理、批处理底层统一的执行引擎,提供对无界表(无边界的源源不断到达的流资料)和有界表(静态历史资料)的优化查询,而向用户提供 Dataset/DataFrame API 对批流资料联合处理,进一步模糊了批流资料处理的边界。
另一方面,Apache Flink 在 2016 年左右进入大众视野,凭借其当时更优的流处理引擎,原生的 Watermark 支援‘Exaclty Once’的资料一致性保证,和批流一体计算等各种场景的支援,成为 Spark 的有力挑战者。无论是使用 Spark 还是 Flink,使用者真正关心的是如何更好地使用资料,更快地挖掘资料中的价值,流资料和静态资料不再是分离的个体,而是一份资料的两种不同表征方式。
然而在实践中,构建一个批流一体的资料平台并不只是计算引擎层的任务。因为在传统解决方案中,近实时的流、事件资料通常采用讯息伫列(例如 RabbitMQ)、实时资料管道(例如 Apache Kafka)储存,而批处理所需要的静态资料通常使用档案系统、物件储存进行储存。这就意味着,一方面,在资料分析过程中,为了保证结果的正确性和实时性,需要对分别储存在两类系统中资料进行联合查询;另一方面,在运维过程中,需要定期将流资料转存到档案 / 物件储存中,通过维持流形式的资料总量在阈值之下来保证讯息伫列、资料管道的效能(因为这类系统的以分割槽为主的架构设计紧耦合了讯息服务和讯息储存,而且多数都太过依赖档案系统,随着资料量的增加,系统性能会急剧下降),但人为的资料搬迁不但会提升系统的运维成本,而且搬迁过程中的资料清洗、读取、载入也是对丛集资源的巨大消耗。
与此同时,从 Mesos 和 YARN 的流行、Docker 的兴起到现在的 Kubernetes 被广泛采用,整个基础架构正在全面地向容器化方向发展,传统紧耦合讯息服务和讯息计算的架构并不能很好地适应容器化的架构。以 Kafka 为例,其以分割槽为中心的架构紧耦合了讯息服务和讯息储存。Kafka 的分割槽与一台或者一组物理机强系结,这带来的问题是在机器失效或丛集扩容中,需要进行昂贵且漫长的分割槽资料重新均衡的过程;其以分割槽为粒度的储存设计也不能很好利用已有的云端储存资源;此外,过于简单的设计导致其为了进行容器化需要解决多租户管理、IO 隔离等方面很多架构上的缺陷。
Pulsar 简介
Apache Pulsar 是一个多租户、高效能的企业级讯息释出订阅系统,最初由 Yahoo 研发, 2018 年 9 月从 Apache 孵化器毕业,成为 Apache 基金会的顶级开源专案。Pulsar 基于释出订阅模式(pub-sub)构建,生产者(producer)释出讯息(message)到主题(topic),消费者可以订阅主题,处理收到的讯息,并在讯息处理完成后传送确认(Ack)。Pulsar 提供了四种订阅型别,它们可以共存在同一个主题上,以订阅名进行区分:
独享(exclusive)订阅——一个订阅名下同时只能有一个消费者。
共享(shared)订阅——可以由多个消费者订阅,每个消费者接收其中一部分讯息。
失效备援(failover)订阅——允许多个消费者连线到同一个主题,但只有一个消费者能够接收讯息。只有在当前消费者发生失效时,其他消费者才开始接收讯息。
键划分(key-shared)订阅(测试版功能)——多个消费者连线到同一主题,相同 Key 总会发送给同一个消费者。
Pulsar 从设计之初就支援多租户(multi-tenancy)的概念,租户(tenant)可以横跨多个丛集(clusters),每个租户都有其认证和鉴权方式,租户也是储存配额、讯息生存时间(TTL)和隔离策略的管理单元。Pulsar 多租户的特性可以在 topic URL 上得到充分体现,其结构是persistent://tenant/namespace/topic。名称空间(namespace)是 Pulsar 中最基本的管理单元,我们可以设定许可权、调整复制选项、管理跨丛集的资料复制、控制讯息的过期时间或执行其他关键任务。
Pulsar 独特架构
Pulsar 和其他讯息系统的最根本区别在于其采用计算和储存分离的分层架构。Pulsar 丛集由两层组成:无状态服务层,它由一组接受和传递讯息的 broker 组成;分散式储存层,它由一组名为 bookies 的 Apache BookKeeper 储存节点组成,具备高可用、强一致、低延时的特点。

和 Kafka 一样,Pulsar 也是基于主题分割槽(Topic partition)的逻辑概念进行主题资料的储存。不同的是,Kafka 的物理储存也是以分割槽为单位,每个 partition 必须作为一个整体(一个目录)被储存在一个 broker 上,而 Pulsar 的每个主题分割槽本质上都是储存在 BookKeeper 上的分散式日志,每个日志又被分成分段(Segment)。每个 Segment 作为 BookKeeper 上的一个 Ledger,均匀分布并存储在多个 bookie 中。储存分层的架构和以 Segment 为中心的分片储存是 Pulsar 的两个关键设计理念。以此为基础为 Pulsar 提供了很多重要的优势:无限制的主题分割槽、储存即时扩充套件,无需资料迁移 、无缝 broker 故障恢复、无缝丛集扩充套件、无缝的储存(Bookie)故障恢复和独立的可扩充套件性。
讯息系统解耦了生产者与消费者,但实际的讯息本质上仍是有结构的,因此生产者和消费者之间需要一种协调机制,达到生产、消费过程中对讯息结构的共识,以达到型别安全的目的。Pulsar 有内建的 Schema 注册方式在讯息系统端提供传输讯息型别约定的方式,客户端可以通过上传 Schema 来约定主题级别的讯息型别资讯,而由 Pulsar 负责讯息的型别检查和有型别讯息的自动序列化、反序列化,从而降低多应用间的讯息解析程式码反复开发、维护的成本。当然,Schema 定义与型别安全是一种可选的机制,并不会给非型别化讯息的释出、消费产生任何效能开销。
在 Spark 中实现对 Pulsar 资料的读写——Spark Pulsar Connector
自 Spark 2.2 版本 Structured Streaming 正式释出,Spark 只保留了 SparkSession 作为主程式入口,你只需编写 DataSet/DataFrame API 程式,以宣告形式对资料的操作,而将具体的查询优化与批流处理执行的细节交由 Spark SQL 引擎进行处理。对于一个数据处理作业,需要定义 DataFrame 的产生、变换和写出三个部分,而将 Pulsar 作为流资料平台与 Spark 进行整合正是要解决如何从 Pulsar 中读取资料(Source)和如何向 Pulsar 写出运算结果(Sink)两个问题。
为了实现以 Pulsar 为源读取批流资料与支援批流资料向 Pulsar 的写入,我们构建了 Spark Pulsar Connector。
对 Structured Streaming 的支援

上图展示了 Structured Streaming(以下简称 SS )的主要元件:
输入和输出——为了提供细粒度的容错,SS 要求输入资料来源(Source)是可重放(replayable)的;为了提供端到端的 Exactly-Once 的语义,需要输出(Sink)支援幂等写出(一条讯息被多次写入与一次写入效果一致,可由 DBMS、KV 系统通过键约束的方式支援)。
API——使用者通过编写 Spark SQL 的 batch API(SQL 或 DataFrame)指定对一个或多个流、表的查询,并定义一个输出表储存所有的输出结果,而引擎内部决定如何将结果增量地写到 Sink 中。为了支援流处理,SS 在原有的 Spark SQL API 上添加了一些界面:
触发器(Trigger)——控制引擎触发流处理执行、在 Sink 中更新结果的频率。
水印机制(Watermark policy)——使用者通过指定字段做 event time,来决定对晚到资料的处理。
有状态算子(Stateful operator)——使用者可以根据 Key 跟踪和更新算子内部的可变状态,完成复杂的业务需求(例如,基于会话的视窗)。
执行层——当收到一个查询时,SS 决定它的增量执行方式,进行优化、并开始执行。SS 有两种可选的执行模型:
Microbatch model(微批处理模式)——预设的执行方式,与 Spark Streaming 的 DStream 类似,将流切成 micro batch,对每个 batch 分别处理。这种模式支援动态负载均衡、故障恢复等机制,适合将吞吐率作为主要效能指标的应用。
Continuous mode(持续模式)——在丛集上启动长时间执行的算子,适合处理较为简单、延迟敏感类应用。
Log 和 State Store —— SS 利用两种持久化储存来提供容错保障:一个 Write-ahead-Log(WAL),记录被成功消费且持久化写出的每个资料来源中的位置;一个大规模的 state store, 储存长期执行的聚集算子内部的状态快照。当故障发生时,SS 会根据快照的位置,通过重放之后的讯息完成流处理状态的恢复。
具体到源代码层面,Source 界面定义了可重放资料来源需要提供的功能。
trait Source {
def schema: StructType
def getOffset: Option[Offset]
def getBatch(start: Option[Offset], end: Offset): DataFrame
def commit(end: Offset): Unit
def stop(): Unit
}
trait Sink {
def addBatch(batchId: Long, data: DataFrame): Unit
}
以 microbatch 执行模式为例:
在每个 microbatch 的最开始,SS 会向 source 询问当前的最新进度(getOffset),并将其持久化到 WAL 中。
随后,source 根据 SS 提供的 start end 偏移量,提供区间范围的资料(getBatch)。
SS 触发计算逻辑的优化和编译,把计算结果写出给 sink(addBatch),这时才触发实际的取资料操作以及计算过程。
在资料完整写出到 sink 后,SS 通知 source 可以废弃资料(commit),并将成功执行的 batchId 写入内部维护的 commitLog 中。
具体到 Pulsar 的 connector 实现中:
在所有批次开始执行前,SS 会呼叫 schema 方法返回讯息的结构资讯,在 schema 方法内部,我们从 Pulsar 的 Schema Registry 提取出所有主题的 Schema,并进行一致性检查。
随后,我们为每个主题分割槽建立一个消费者,按照 (start, end] 返回主题分割槽中的资料。
当收到 SS 的 commit 通知时,通过 topics 中的 resetCursor 向 Pulsar 标志讯息消费的完成。Sink 中构建的生产者则将 addBatch 中获取的实际资料以讯息形式追加写入相应的主题中。

对批处理作业的支援
在某个时间点执行的批作业,可以看作是对 Pulsar 平台中的流资料在一个时间点的快照进行的资料分析。Spark 对历史资料的查询是以 Relation 为单位,Spark Pulsar Connector 提供 createRelation 方法的实现根据使用者指定的多个主题分割槽构建表,并返回包含 Schema 资讯的 DataSet。在查询计划阶段,Connector 的功能分成两步:首先,根据使用者提供的一个或多个主题,在 Pulsar Schema Registry 中查询主题 Schema,并检查多个主题 Schema 的一致性;其次,将使用者指定的所有主题分割槽进行任务划分(Partition),得到的分片即是 Spark source task 的执行粒度。
Pulsar 提供了两层的界面对其中的资料进行访问,基于主题分割槽的 Consumer/Reader 界面,以传统讯息接收为语义的顺序资料读取;Segment 级的读界面,提供对 Segment 资料的直接读取。因此,相应地从 Pulsar 读资料执行批作业可以分成两种粒度(即读取资料的并行度)进行:以主题分割槽为粒度(每个主题分割槽作为一个分片);以 Segment 为粒度(将一个主题分割槽的多个 Segment 组织成一个分片,因此一个主题分割槽会有多个对应的分片)。你可以按照批作业的并行度需求和可分配计算资源选择合适的讯息读取的并行粒度。另一方面,将批作业的执行储存到 Pulsar 也很直观,你只需指定写入的主题和讯息路由规则(RoundRobin 或者按 Key 划分),在 Sink task 中建立的每个生产者会将待写出的讯息送至对应的主题分割槽。
如何使用 Spark Pulsar Connector
根据一个或多个主题建立流处理 Source。
val df = spark
.readStream
.format("pulsar")
.option("service.url", "pulsar://localhost:6650")
.option("admin.url", "http://localhost:8080")
.option("topicsPattern", "topic.*") // Subscribe to a pattern
// .option("topics", "topic1,topic2") // Subscribe to multiple topics
// .option("topic", "topic1"). //subscribe to a single topic
.option("startingOffsets", startingOffsets)
.load()
df.selectExpr("CAST(__key AS STRING)", "CAST(value AS STRING)")
.as[(String, String)]
构建批处理 Source。
val df = spark
.read
.format("pulsar")
.option("service.url", "pulsar://localhost:6650")
.option("admin.url", "http://localhost:8080")
.option("topicsPattern", "topic.*")
.option("startingOffsets", "earliest")
.option("endingOffsets", "latest")
.load()
df.selectExpr("CAST(__key AS STRING)", "CAST(value AS STRING)")
.as[(String, String)]
使用资料中本身的 topic 字段向多个主题进行持续 Sink。
val df = spark
.read
.format("pulsar")
.option("service.url", "pulsar://localhost:6650")
.option("admin.url", "http://localhost:8080")
.option("topicsPattern", "topic.*")
.option("startingOffsets", "earliest")
.option("endingOffsets", "latest")
.load()
df.selectExpr("CAST(__key AS STRING)", "CAST(value AS STRING)")
.as[(String, String)]
将批处理结果写回 Pulsar。
df.selectExpr("CAST(__key AS STRING)", "CAST(value AS STRING)")
.write
.format("pulsar")
.option("service.url", "pulsar://localhost:6650")
.option("topic", "topic1")
.save()

注意
由于 Spark Pulsar Connector 支援结构化讯息的消费和写入,为了避免讯息负载中字段和讯息元资料(event time、publish time、key 和 messageId)的潜在命名冲突,讯息元资料字段在 Spark schema 中以双下划线做为字首(例如,__eventTime)。
参考资料
Structured Streaming: A Declarative API for Real-Time Applications in Apache Spark
Structured Streaming 源代码解析系列
Pulsar 官网
从讯息系统到资料平台