Data Engineer 面试题 2026:20 道 SQL、Spark、Kafka、数据建模与系统设计高频题
Data Engineer Interview Questions数据工程师面试题大数据面试题DataEngineeringSQLPythonSparkKafkaETLsystemdesign

Data Engineer 面试题 2026:20 道 SQL、Spark、Kafka、数据建模与系统设计高频题

20 道 Data Engineer 面试题 2026:SQL 窗口函数、Spark 优化、Kafka、数据建模、CDC、数据回填与管道系统设计,附参考答案与追问。

Sam · · 14 分钟阅读

快速答案:Data Engineer 面试题怎么准备

Data Engineer interview questions 通常不是单纯问工具定义,而是考你能不能把 SQL、Python、Spark、Kafka、数据建模和 pipeline reliability 放到真实业务场景里回答。准备数据工程师面试题时,优先级应该是:先练 SQL 和数据建模,再练 Spark/Kafka 的大数据面试题,最后把实时管道、数据质量、回填、幂等性和系统设计串起来。

如果你时间有限,先按下面 4 类题准备:

题型面试官想看什么本文对应题目
SQL / Python coding数据处理基本功、窗口函数、边界条件13、14、20
Spark / Kafka / Big Data大数据框架原理、分区、shuffle、延迟数据4、5、6、15、17
Data Modeling / Warehouse事实表、维度表、SCD、订单模型实战1、7、8、19
Pipeline System Design吞吐估算、容错、幂等性、数据质量、监控3、9、10、11、12、18
Architecture / Trade-offsLambda vs Kappa、湖仓取舍、技术选型2、16

如果你想按面试轮次系统准备,可以先看 Data Engineer 面试准备完全指南,再用本文刷高频题。需要一对一 mock 的话,可以看 Data Engineer 面试辅导

数据工程面试很独特——它们融合了编码能力分布式系统知识实际架构经验。与偏重算法的 SDE 面试不同,数据工程面试考察的是你是否能设计和构建现实世界的数据基础设施。

以下是 20 道最常见的题目,附带详细答案和准备策略。

数据工程核心概念

1. 数据仓库和数据湖有什么区别?

数据仓库:

  • 仅结构化数据
  • Schema-on-write(数据在存储前验证)
  • 为 SQL 分析优化
  • 举例:Snowflake、Redshift、BigQuery
  • 使用场景:BI 仪表盘、报表、结构化分析

数据湖:

  • 任何格式的原始数据(结构化、半结构化、非结构化)
  • Schema-on-read(查询时应用 schema)
  • 存储文件(Parquet、Avro、JSON、CSV、图片、日志)
  • 举例:AWS S3、Azure Data Lake、GCP Cloud Storage
  • 使用场景:ML 训练数据、日志分析、探索性数据科学

湖仓一体(Lakehouse,现代方法): 两者结合——在湖中存储原始数据,同时加入 ACID 事务和查询优化。举例:Delta Lake、Apache Iceberg、Apache Hudi。

面试官为什么问: 他们想看看你是否理解何时使用每种架构,并能讨论各自的权衡。

2. 解释 ETL 与 ELT 的区别

ETL(提取 → 转换 → 加载):

  • 在加载到数据仓库之前转换数据
  • 传统方法,用于数据仓库计算成本较高的场景
  • 工具举例:Informatica、传统 SSIS

ELT(提取 → 加载 → 转换):

  • 先加载原始数据,然后在数据仓库内转换
  • 现代方法,得益于廉价的云存储和强大的数据仓库计算能力
  • 工具举例:dbt(转换步骤)、Fivetran/Singer(提取/加载)

关键洞察: 行业已经大幅转向 ELT。云数据仓库(Snowflake、BigQuery)使得转换成本很低。主要的权衡是治理——ELT 意味着原始(可能杂乱的)数据在你的数据仓库中停留更久。

3. 数据管道中的幂等性是什么,为什么重要?

幂等性意味着多次运行同一个管道与运行一次产生相同的结果。这很关键,因为:

  • 管道会失败,需要重试
  • 回填历史数据会重新处理已加载的数据
  • 在分布式系统中保证精确一次处理很难

如何实现幂等性:

  • 使用唯一键(upserts 而非 inserts)
  • 按时间窗口分区并覆盖
  • 在数据仓库中使用 merge 操作
  • 在暂存表中跟踪已处理的事件 ID

大数据框架

4. Apache Spark 如何处理数据分区?

Spark 在内存中(计算期间)和磁盘上(读写时)对数据进行分区。关键概念:

  • 默认分区: 集群中的核心数(或 RDD 操作的 200 个分区)
  • Repartition: 将所有数据通过网络重新洗牌——代价高昂,谨慎使用
  • Coalesce: 减少分区而不进行完全 shuffle——缩小规模时效率高
  • Partition by: 写入时按列将数据分组到文件(例如 partitionBy("date")

常见面试追问: “如果你遇到数据倾斜怎么办?”

回答:一个分区比其他分区接收更多的数据,导致某些任务运行时间更长。解决方案:

  • 加盐(给 key 添加随机前缀)
  • 自定义分区器
  • 对小表使用广播连接

5. 解释 Kafka 与其他消息队列(RabbitMQ、SQS)的区别

特性KafkaRabbitMQAWS SQS
模式发布/订阅(事件流)消息代理(队列)简单队列
持久化是(可配置保留期)可选
有序性按分区有序按队列有序尽力而为(FIFO 队列)
回放是(保留消息)
吞吐量百万/秒数千/秒数十万/秒
使用场景事件流、管道任务分发微服务解耦

关键洞察: Kafka 专为高吞吐量事件流设计,支持保留和回放。RabbitMQ 在消息路由和复杂模式方面表现优异。SQS 是最简单的,用于基本的解耦场景。

6. 如何在流式管道中处理延迟到达的数据?

延迟数据在实时系统中是不可避免的。策略如下:

  1. 水位线(Watermarks): 定义一个阈值(例如 5 分钟)——在水位线之后到达的数据被视为延迟
  2. 允许延迟的会话窗口: 在窗口中处理数据,但允许延迟事件更新结果
  3. 侧输出(Side outputs): 将延迟数据路由到单独的流进行分析
  4. 状态化处理: 在内存/磁盘上保持状态以更新之前的结果
  5. 定期批处理对账: 运行每日批处理作业来纠正任何流式处理的不准确之处

现实案例: 在点击流管道中,网络不佳的移动用户可能延迟 10-30 秒到达事件。5 分钟水位线处理了 99.9% 的情况,每日批处理作业修复剩余部分。

数据建模

7. 什么是维度建模?解释星型和雪花模式的区别。

维度建模(Kimball 方法论)使用以下方式组织数据以支持分析查询:

  • 事实表: 定量、可测量的数据(销售额、点击量、交易)
  • 维度表: 描述性上下文(产品、客户、时间、地点)

星型模式(Star Schema):

    Dim_Time    Dim_Customer
       \         /
        \       /
         Fact_Sales
        /       \
       /         \
  Dim_Product  Dim_Store
  • 事实表直接连接所有维度
  • 查询更简单,性能更快
  • 维度表中有一些数据冗余

雪花模式(Snowflake Schema):

  • 维度表规范化(拆分为子维度)
  • 冗余更少,查询更复杂
  • 在现代云数据仓库中较少使用(存储成本低)

面试提示: 大多数从业者偏好星型模式,因其简单且查询性能好。雪花模式在理论上更优雅,但实际上更慢。

8. SCD(缓慢变化维)类型 1 和类型 2 的区别?

类型 1(覆盖): 用新值替换旧值。无历史记录。

  • 适用场景:历史记录不重要(例如,修正拼写错误)

类型 2(新增行): 创建带有新生效日期的新行。保留历史记录。

  • 适用场景:需要跟踪随时间的变化(例如,客户地址变更)
  • 需要字段:effective_dateexpiry_dateis_current

类型 3(新增列): 添加”旧值”列。有限的历史记录。

  • 适用场景:只需要上一状态,不需要完整历史

真实面试题: “你的分析数据库中,一个客户从免费版升级到专业版,你会如何处理?”

回答:类型 2 SCD——你需要知道每笔交易发生时客户属于哪个层级。

管道设计

9. 为网约车应用设计一个实时分析管道。

需求: 近实时跟踪乘车请求、匹配、完成和收入。

架构:

Mobile App → API Gateway → Kafka (events) → Spark Streaming / Flink

                                              Aggregation (per driver, per zone)

                                              Real-time Dashboard (Redis/Memcached)

                                              Data Warehouse (historical analysis)

关键组件:

  1. 事件采集: Kafka 主题用于 ride_requestedride_matchedride_completed
  2. 流处理: 连接流(将请求与完成匹配)、按司机/区域聚合
  3. 状态管理: 在状态存储(Redis)中跟踪活跃的乘车
  4. 输出: 实时指标到仪表盘 + 批量写入数据仓库进行历史分析
  5. 监控: 在管道延迟、数据质量问题、schema 变更时告警

预期追问:

  • 如何实现精确一次处理?
  • Kafka 宕机时怎么办?
  • 如何处理 schema 演进?

10. 如何确保管道中的数据质量?

采集阶段:

  • Schema 验证(检查列类型、必填字段)
  • 关键列的空值检查
  • 范围验证(例如,时间戳必须在过去)

处理阶段:

  • 行数检查(输入 vs 输出)
  • 重复检测
  • 参照完整性检查

输出阶段:

  • 与之前运行对比的摘要统计
  • 关键指标的异常检测
  • 新鲜度监控(数据到达时间)

工具: Great Expectations、dbt tests、Deequ、自定义验证作业。

系统设计

11. 如何设计一个每天处理 100 亿事件的管道?

关键考量:

  1. 分区: 按日期 + 小时 + 区域对事件分区,实现并行处理
  2. 存储: 将原始数据以 Parquet 格式(压缩、列式)存入 S3/GCS
  3. 处理: Spark Structured Streaming 或 Flink 用于实时处理,Spark 批处理用于每日作业
  4. 编排: Airflow/Dagster 协调依赖关系
  5. 扩展: 基于队列深度自动扩展集群
  6. 成本优化: 容错批处理作业使用 Spot 实例,关键路径使用预留实例

数学计算: 100 亿事件 × 约 1KB/事件 = 约 10TB/天。Parquet 10:1 压缩后,约 1TB 存储数据。100 个执行器的 Spark 集群约 15 分钟处理完毕。

12. 什么是 CDC(变更数据捕获)管道,它是如何工作的?

CDC 实时跟踪数据库变更,无需轮询。

工作原理:

  1. 读取数据库的事务日志(PostgreSQL 的 WAL、MySQL 的 binlog)
  2. 解析 INSERT/UPDATE/DELETE 操作
  3. 将变更流式传输到 Kafka
  4. 消费变更以更新下游系统

工具: Debezium(开源)、AWS DMS、Fivetran(SaaS)

相比轮询的优势:

  • 不增加数据库负载(读取日志成本低)
  • 亚秒级延迟
  • 精确一次捕获每次变更

数据工程师编码题

13. 编写 Python 函数,按 event_id 去重,保留最新的时间戳

from collections import defaultdict

def deduplicate_events(events: list[dict]) -> list[dict]:
    """
    events: list of dicts with 'event_id' and 'timestamp' keys
    Returns deduplicated list keeping the latest event for each event_id
    """
    latest = {}
    for event in events:
        eid = event['event_id']
        if eid not in latest or event['timestamp'] > latest[eid]['timestamp']:
            latest[eid] = event
    return list(latest.values())

追问: 如果数据集放不下内存怎么办? 回答:按 event_id, timestamp DESC 排序,然后使用流式方法(只保留最后看到的 event_id)。

14. SQL:找出连续 7 天活跃的用户

WITH active_days AS (
  SELECT DISTINCT user_id, DATE(activity_date) as activity_date
  FROM user_activity
),
streaks AS (
  SELECT
    user_id,
    activity_date,
    DATE_SUB(activity_date, ROW_NUMBER() OVER (
      PARTITION BY user_id ORDER BY activity_date
    )) as streak_group
  FROM active_days
)
SELECT user_id, COUNT(*) as streak_length
FROM streaks
GROUP BY user_id, streak_group
HAVING COUNT(*) >= 7;

关键洞察: DATE_SUB 技巧将连续日期转换为同一组。如果日期连续,减去行号会得到相同的结果。

15. 如何优化一个慢速的 Spark 作业?

常见优化技巧:

  1. 减少 shuffle: 对小表使用 broadcast(),用 mapJoin 替代 join
  2. 合理分区: 按常用过滤列分区,避免过多小文件
  3. 缓存中间结果: 多次使用的 dataframe 使用 cache()persist()
  4. 尽早过滤: 在 join 之前下推过滤条件
  5. 使用合适的文件格式: Parquet 配合列裁剪
  6. 处理数据倾斜: 加盐、自定义分区器、单独处理热 key
  7. 用 Spark UI 监控: 识别慢速 stage 和长尾任务

进阶高频题

16. Lambda 架构和 Kappa 架构有什么区别?

Lambda 架构: 三层设计——

  1. Batch layer: 用全量数据计算完整、正确的结果(Spark/Hadoop 批处理),可反复重算纠错
  2. Speed layer: 用流处理(Flink/Spark Streaming)快速产出低延迟的近似结果
  3. Serving layer: 合并两层结果对外服务

优点: 容错强、可修正历史错误。缺点: 同一套逻辑要写两遍(批 + 流),维护成本高,两层结果在切换时可能短暂不一致。

Kappa 架构: 去掉 batch layer,一切皆流。实时计算 + 需要重算时回放 Kafka/日志数据重新跑一遍流。

优点: 只维护一条代码路径,架构简单。缺点: 强依赖消息队列的保留时长和回放能力,长周期历史回放在成本和状态管理上都比较重。

面试回答建议: 不要死背架构名称。实际生产里更常见的是”Kappa 思路 + 批处理兜底”:流式管道负责实时,定期批作业负责纠错和对账,用 Delta Lake/Iceberg 这类湖仓格式统一存储,避免真正的双代码路径。

17. Parquet、ORC、CSV 文件格式怎么选?

维度CSVParquetORC
格式行式列式列式
压缩率高(通常 5-10x)
列裁剪不支持支持支持
Predicate pushdown不支持支持支持
Schema无(需推断)内嵌内嵌
典型生态数据交换Spark / BigQuery / 大多数湖仓Hive

选择建议:

  • 分析场景默认选 Parquet——列裁剪 + 压缩让 IO 成本大幅下降
  • CSV 只用于小规模数据交换或日志原始落盘,不要作为分析层格式
  • 需要 ACID、upsert、time travel 时,在 Parquet 文件之上加 Delta Lake / Iceberg 表格式
  • 追问预警: 小文件问题——流式写入会产生大量小文件,需要定期 compaction(coalesce/rewrite 作业)

18. 如何设计一个数据回填(backfill)策略?

场景:上游 schema 变了、聚合逻辑发现 bug、要补算过去 N 个月的历史数据。

标准步骤:

  1. 圈定范围: 明确哪些表、哪个时间区间、下游影响面
  2. 隔离执行: 用独立调度(不走日常调度),先写入 staging 表,不直接覆盖生产
  3. 保证幂等: 按分区 overwrite(如 INSERT OVERWRITE PARTITION (dt=...)),重复跑结果一致
  4. 分批控制: 按天/小时分批跑,限制并发,避免把生产集群资源打满
  5. 验证: 行数核对、关键指标对比(新旧逻辑 diff)、跑 dbt tests / Great Expectations
  6. 切换: 验证通过后原子切换下游指向,旧数据保留一段时间可回滚
  7. 沟通: 提前通知下游团队回填窗口内数据可能波动

现实案例: 发现某聚合指标算错需要回填 90 天——先回填 3 天做 pilot,确认耗时和成本在预算内,再分批跑完剩余部分,最后用一周新旧数据对比确认指标口径一致。

19. 数据建模实战:为电商订单系统设计数仓 schema

第一步——先定义 grain(粒度): 一行事实表代表什么?订单级(一行 = 一个订单)还是订单行级(一行 = 一个商品项)?先回答这个,后面才不会乱。

事实表:

  • fact_orders:order_id、user_id、product_id、order_time、quantity、amount、discount、status
  • 如果要跟踪订单状态流转(待支付 → 已支付 → 已发货 → 已完成/退款),用累积快照事实表(accumulating snapshot),每行记录各状态发生的时间戳

维度表:

  • dim_user(用户)、dim_product(商品)、dim_time(时间)、dim_channel(渠道)
  • 缓慢变化:用户地址变更 → SCD Type 2;商品价格变化 → Type 2 或 Type 3

聚合层:

  • 预聚合 daily_sales_rollup(日期 × 品类 × 渠道粒度的 GMV、订单数、退款率),BI 高频查询直接打聚合表,避免每次扫全量事实表

面试加分点: 主动说明如何支撑典型查询(“按天看各品类 GMV 和客单价”),并提到聚合表刷新频率与事实表的一致性(T+1 vs 近实时)。

20. SQL:找出每个类目下销售额 Top 3 的商品(Top N per group)

WITH ranked AS (
  SELECT
    category,
    product_id,
    SUM(amount) AS total_sales,
    DENSE_RANK() OVER (
      PARTITION BY category ORDER BY SUM(amount) DESC
    ) AS rnk
  FROM orders
  WHERE order_date >= DATE_SUB(CURRENT_DATE, 30)
  GROUP BY category, product_id
)
SELECT category, product_id, total_sales
FROM ranked
WHERE rnk <= 3
ORDER BY category, total_sales DESC;

关键点:

  • 窗口函数里不能直接引用聚合别名,所以先 GROUP BY 再开窗(包一层 CTE)
  • DENSE_RANK() 而不是 ROW_NUMBER():并列销售额时不会漏掉商品
  • 必考追问: ROW_NUMBER / RANK / DENSE_RANK 有什么区别?——分别回答”强制唯一序号 / 并列跳号 / 并列不跳号”,并举例说明三种结果差异

高频题地图(按框架)

SQL 高频题型

题型典型题目本文对应
窗口函数Top N per group、running total、排名Q20
连续日期(streak)连续 N 天活跃Q14
去重保留最新一条、行级去重Q13
Retention / Cohort次日/7 日留存、cohort 分析DataLemur 练习
自连接同比/环比、与前一天对比DataLemur 练习
聚合 + 子查询类目级指标、分位数

准备建议:DataLemur / StrataScratch 刷 40-60 题即可覆盖 90% 场景,优先上面前 4 类。每题自己写完 SQL 再对答案,并能口头解释执行计划(哪个操作触发 shuffle、窗口函数为什么快)。

Spark / Kafka / Airflow 高频题型

框架高频考察方向本文对应
Spark分区、shuffle、数据倾斜、优化 checklist、文件格式选择Q4、Q15、Q17
Kafka与其他 MQ 对比、分区与有序性、消费者组、回放、watermark 与延迟数据Q5、Q6
Airflow / 编排DAG 设计、backfill、幂等任务、SLA 监控、动态 DAG见下

Airflow 高频题速查(DE 面试最容易漏的一块):

  • DAG 怎么设计:任务粒度怎么拆、依赖怎么表达、为什么避免一条超长的链(调试和重跑成本)
  • Backfill / 历史补数:怎么分批跑、怎么限制并发、怎么验证结果一致(与本文 Q18 配套回答)
  • 任务幂等:retry-safe 的写法(overwrite 而非 append)、partial failure 怎么处理
  • SLA 与监控:上游延迟怎么告警、数据新鲜度怎么定义、on-call 响应流程
  • 动态 DAG vs 静态 DAG:基于元数据/配置生成 DAG 的取舍(灵活 vs 难审计)

延伸阅读:DE 案例:数据管道编排系统设计(调度、依赖、失败重试的完整设计)

DE System Design 题型地图

DE 的系统设计不是 URL shortener,而是下面 6 类。回答时从数据量、延迟、可靠性、schema、回填和监控切入,最后再落到工具:

题型核心考察点本文题目深度案例
实时分析管道事件接入、流处理、聚合、dashboardQ9实时数据湖架构案例
数据回填 / 历史重算圈定范围、隔离、幂等、分批、验证Q18数据管道编排设计案例
数据质量平台校验规则、异常检测、新鲜度监控Q10Data Engineer 准备指南
湖仓 / 数仓建模Parquet + Iceberg/Delta、分层架构、SCDQ1、Q7、Q19实时数据湖架构案例
CDC 实时同步WAL/binlog、Debezium、exactly-onceQ12Streaming DE 辅导案例
用户行为 / 画像平台事件模型、画像更新、多租户隔离多租户数据隔离案例Data Analyst 转 DE 案例

准备策略

第 1-2 周: SQL 进阶(窗口函数、CTE、自连接)——在 DataLemur 上练习 第 3-4 周: Spark 基础 + 动手编码练习 第 5-6 周: 系统设计——练习在白板上设计管道 第 7-8 周: Kafka、消息队列、流式处理概念 第 9 周及以后: 模拟面试,聚焦管道设计场景

资源推荐

  • 书籍: Martin Kleppmann 的《Designing Data-Intensive Applications》(必读)
  • SQL 练习: DataLemur、StrataScratch
  • Spark: Databricks 免费社区版
  • 系统设计: Alex Xu 的《System Design Interview》(第 2 卷涵盖数据系统)

总结

数据工程面试考察的是你思考规模可靠性权衡的能力。你不需要掌握每一个工具——你需要理解它们背后的原理,并能够设计在生产环境中工作的系统。

FAQ

Data Engineer interview questions 最常考哪些主题?

最常见的是 SQL、Python 数据处理、Spark 分区和 shuffle、Kafka 事件流、ETL/ELT、数据建模、CDC、数据质量和 pipeline system design。Senior Data Engineer 还会被追问成本优化、backfill、exactly-once 语义、schema evolution 和 observability。

大数据面试题是不是只要会 Spark 就够了?

不够。Spark 是核心,但面试官更关心你是否理解分区、shuffle、数据倾斜、文件格式、watermark、延迟数据、checkpoint、重试和幂等性。好的回答要把 Spark、Kafka、Airflow/dbt、warehouse/lakehouse 和监控串成完整的数据平台。

数据工程师面试需要准备系统设计吗?

需要。Data Engineer 的系统设计通常不是 URL shortener,而是数据管道、实时分析、CDC、data lake、warehouse modeling、feature pipeline 或数据质量平台。回答时要从数据量、延迟、可靠性、schema、回填和监控开始,而不是只画工具链。

Data Engineer 面试 SQL 和编码各占多少比重?

SQL 是基础题,几乎每轮都会出现(窗口函数、连续活跃、Top N per group、retention 是四大高频)。Python/编码比重比 SDE 低,主要考数据处理(去重、聚合、流式处理)而不是算法。真正拉开差距的是数据建模和管道设计——把 SQL 当基本功练熟,把时间花在”设计一个管道”这种开放题上。

数据工程师和数据分析师的面试题差别在哪?

分析师题偏 SQL + 统计学 + 业务指标解读(A/B test、cohort、归因);工程师题偏管道工程 + 分布式系统(Spark 分区、Kafka 语义、幂等性、backfill、数据质量)。SQL 两边都考,但工程师的 SQL 更侧重窗口函数和大规模查询优化。如果是转型岗位,两边都要准备,但按目标 JD 的侧重点分配时间。

时间有限,Data Engineer 面试题先刷什么?

优先级:SQL 窗口函数和经典题(第 1、2 周)→ 数据建模(star schema、SCD、订单模型)→ Spark 分区/shuffle/优化 + Kafka 核心概念 → 管道系统设计(幂等、数据质量、backfill、10B 事件)。工具名词不用背全,但要能解释分区、shuffle、watermark、exactly-once 背后的原理。

不是必须。Spark Structured Streaming 是覆盖面最广的答案,会 Spark 流处理就足够应对大多数公司。Flink 在强流处理场景(金融、广告实时归因)更常见,如果你目标公司技术栈是 Flink,再补 state、checkpoint、watermark 这些概念。面试里说”我们生产用 Spark Streaming,Flink 的 checkpoint 机制我了解过”比硬装精通安全得多。


推荐阅读


准备好和经历过真实面试流程的数据工程师一起练习了吗?查看 Data Engineer 面试辅导,我们会评估你的 SQL、建模和系统设计水平,制定针对性准备计划;也可以直接联系我们

S

关于作者

Sam 是 Interview Coach Pro 的技术面试教练,长期辅导在美国求职的中文候选人准备 SDE、System Design、Behavioral、Data Engineer 和 ML Engineer 面试。

本文基于匿名面试复盘、公开岗位要求和一对一辅导中的高频问题整理,发布前会检查内容结构、术语准确性和可操作性。你也可以查看我们的 辅导团队辅导方法

相关面试辅导

如果你正在准备类似面试,可以直接从下面的专项辅导开始。

准备好拿下下一次面试了吗?

获取针对你的目标岗位和公司的个性化辅导方案。

联系我们