数据岗位面试题更新 2026-08-05

在使用 PySpark 进行实时流式数据处理时,通常需要与 Kafka 集成。请阐述 PySpark Structured Streaming 对接 Kafka 的典型实现流程、关键配置项以及消费数据的处理方式。

数据风险判断系统设计技术原理PySpark

考察说明

考察对 PySpark Structured Streaming 与 Kafka 集成机制及实践细节的理解。

回答思路

  1. 【回答框架 1】Structured Streaming 通过 DataStreamReader 的 format 设为 kafka 来读取 Kafka 数据,需指定 kafka.bootstrap.servers 和 subscribe 或 subscribePattern 等参数,返回的 DataFrame 包含 key、value、topic、partition、offset 等列。
  2. 【回答框架 2】处理时通常将 value 列按业务格式(如 JSON)解析,使用 from_json 配合 schema 转换为结构化数据,再进行窗口、聚合等操作,最后用 DataStreamWriter 输出到外部系统,支持 Kafka、文件、控制台等多种 sink。
  3. 【回答框架 3】确保精确一次语义需在输出时设置 checkpointLocation 并启用相应的事务性 sink,或通过幂等写入实现;需理解处理延迟与吞吐量之间的权衡,合理设置触发间隔和并行度。
  4. 【回答框架 4】密钥等敏感信息应通过配置或环境变量安全传递,并注意 Kafka 版本与 PySpark 版本兼容性。
  5. 【关键点 1】使用 PySpark Structured Streaming 读取 Kafka 需指定 format 为 kafka,并配置 bootstrap.servers 和 subscribe 参数。
  6. 【关键点 2】value 字段通常是二进制的,需要反序列化并解析为结构化数据以便处理。
  7. 【关键点 3】通过 checkpoint 支持故障恢复,可实现至少一次或精确一次语义。
  8. 【关键点 4】输出到 Kafka 或其他系统时,可通过写入配置控制语义和数据一致性。
  9. 【易错点 1】忽略 checkpoint 设置可能导致恢复时数据丢失或重复。
  10. 【易错点 2】直接对 value 列进行复杂解析而未定义 schema 会导致性能问题。
  11. 【易错点 3】未考虑反压和背压机制可能造成消费积压或资源浪费。