[toc]
性能优化
大表join怎么优化
- 大表过滤前置,减少数据量
- 合理分区,按时分区裁剪
- 小表广播join
- 倾斜key单独打散
- 开启MapJoin
- 避免笛卡尔积
- 拆分复杂SQL,分步执行。
[toc]
[toc]
离线数仓:T+1 / 小时级,批量调度、跑历史数据,用于报表、复盘、画像;
实时数仓:秒 / 分钟级,CDC 增量流式处理,用于大屏、实时指标、实时风控。
| 维度 | 离线数仓 | 实时数仓 |
|---|---|---|
| 延迟 | 小时级、T+1 | 秒 / 分钟级 |
| 数据方式 | 批量全量 / 定时增量 | CDC 实时增量 |
| 核心引擎 | Hive、Spark | Flink、Kafka |
| 存储介质 | HDFS、Hive | Kafka、Doris、ClickHouse |
| 数据新鲜度 | 低 | 高 |
| 计算模式 | 批处理 | 流处理 |
| 资源开销 | 低、夜间调度 | 常驻集群、资源一直占用 |
| 业务场景 | 日报、月报、复盘、画像、离线 BI | 实时大屏、实时风控、实时订单、实时推荐 |
| 数据精度 | 高、可反复重跑修正 | 准实时、轻微迟到容忍 |
| 运维难度 | 简单、成熟 | 稍复杂、要调优状态 / 反压 / 倾斜 |
写 Hive/Spark SQL,定时调度,跑批、重跑方便,逻辑简单稳定。
Flink SQL + CDC + Kafka,要考虑:
水位线、迟到数据、状态管理、数据倾斜、反压、Exactly-Once、幂等去重。
离线数仓以 Hive+Spark 为核心,采用定时批量抽取、T+1 或小时级调度,适合日报月报、用户画像、业务复盘,数据精准可重跑;
实时数仓基于 Flink+Kafka+CDC 架构,通过监听 MySQL Binlog 做增量同步,秒级延迟,分层实时加工,落地到 Doris/ClickHouse,支撑实时大屏、风控、订单实时指标;
现在企业主流都是离线 + 实时双仓架构,实时看当下、离线做复盘和校准。
用 Hudi/Paimon 做数据湖,一套数据支持:
Flink 实时消费
Spark 离线分析
实现 流批一体、实时离线统一。
[toc]
| 维度 | Apache Doris | ClickHouse |
|---|---|---|
| 数据模型 | 3 种模型:Duplicate (明细)、Aggregate (聚合)、Unique (主键更新) | MergeTree 家族:ReplacingMergeTree、SummingMergeTree 等,模型单一 |
| 事务 / 更新 | 强一致事务,支持同步 UPSERT/DELETE,主键更新性能比 CK 快18–34 倍VeloDB | 无完整事务,更新异步(后台 Merge),存在读写不一致,主键更新极慢 |
| Schema 变更 | 动态调整,支持在线 DDL,成本低 | 变更成本高,常需重建表 |
单表查询:ClickHouse 略优(极致压缩 + 向量化);Doris 接近,ClickBench 互有胜负。
多表 Join :Doris 碾压级优势
并发能力:Doris 支持千级 QPS(BI / 报表);ClickHouse 并发低,高负载易引发 Merge 风暴Apache Doris。
实时更新查询:Doris(Unique Key)比 CK(ReplacingMergeTree)快2.5–4.6 倍(ClickBench)VeloDB。



Doris:
ClickHouse:
Doris:
ClickHouse:
[toc]
| 工具 | 定位 | 性能参考(100MB CSV) | 内存占用 | 关键特点 |
|---|---|---|---|---|
| DataProfiler | 探查与发现 | 2.1s-5 | 45MB-5 | 自动检测格式/PII/统计量 |
| Great Expectations | 验证与测试 | 12.1s--5 | 290MB--5 | 强规则引擎、可创建断言 |
| pandas describe() | 单机探索 | 8.4s-5 | 380MB-5 | 快速简要统计,无法处理超大文件 |
| deequ (Spark) | 分布式验证 | 15.3s-5 | 1.2GB-5 | 依赖 Spark 生态,全量扫描 |
它的核心价值不在于“验证”数据是否满足预期,而在于**“探索和发现”数据中有什么,尤其在自动识别敏感数据方面非常出色**-7。
下面我们来看看它的几项关键能力:
DataProfiler 提供了一个简洁的 API,其内部架构主要围绕以下几个组件构建:
Data 类自动识别并加载多种格式的数据-1。Profiler 类是核心组件,负责协调分析任务-1。下面将 DataProfiler 与其他常用工具进行对比,方便你评估选择:
| 工具 | 定位 | 性能参考(100MB CSV) | 内存占用 | 关键特点 |
|---|---|---|---|---|
| DataProfiler | 探查与发现 | 2.1s-5 | 45MB-5 | 自动检测格式/PII/统计量 |
| Great Expectations | 验证与测试 | 12.1s--5 | 290MB--5 | 强规则引擎、可创建断言 |
| pandas describe() | 单机探索 | 8.4s-5 | 380MB-5 | 快速简要统计,无法处理超大文件 |
| deequ (Spark) | 分布式验证 | 15.3s-5 | 1.2GB-5 | 依赖 Spark 生态,全量扫描 |
DataProfiler 在特定场景下很有优势,但同时也有一些局限性需要注意:
综合来看,DataProfiler 是数据工程师和数据科学家在数据探索、数据发现、敏感数据识别和数据漂移检测等场景下的得力助手。
它通常用于数据管道的以下阶段:
总的来说,DataProfiler 适合在数据生命周期的前期使用,也就是“先探查,后开发”的思路,这一点与其他工具形成了很好的互补。
它通常被用在数据处理管道的前期,也就是在正式开发和测试开始前,先对数据做一个全面的摸底-6。
| 应用阶段 | 核心工作 | DataProfiler 如何发挥作用 |
|---|---|---|
| 1. 🕵️ 阶段一:数据探索与理解 | 对新数据源进行初步摸底,了解其“样子”、基本特征和潜在风险。 | 自动化探索:自动生成数据概要、统计信息和数据模式(schema),帮助快速掌握数据概况-1。 敏感数据发现:利用深度学习模型自动识别CSV、Parquet等格式中的信用卡号、社保号等PII/NPI信息--11。 |
| 2. 🔧 阶段二:开发过程中的数据集成与ETL | 在数据处理初期,保证数据质量,避免“脏数据”污染下游。 | 质量守门员:在ETL任务开始前执行探查,检测数据缺失、类型变化等问题,可配置为任务前置检查-4。 持续监控助手:在开发迭代中,通过对比新旧Profile报告,及时发现代码变更导致的数据不一致-4。 |
| 3. 📈 阶段三:持续的数据监控与数据漂移检测 | 监控生产环境中的数据质量,及早发现因上游变更导致的异常。 | 生成Profile快照:定期(如按日/周)为生产数据生成Profile文件并保存-4。 比对发现异常:对比当前Profile与历史基线,关注核心统计指标(如行数、均值、空值比例)的显著波动,以发现数据漂移。 |
| 4. 🔒 阶段四:数据合规与隐私保护 | 主动发现和定位数据资产中的隐私风险,满足合规要求。 | 主动扫描风险:持续地对公司数据资产进行合规扫描,主动发现是否含有未被声明的PII数据-。 量化隐私风险:在共享数据给团队或合作伙伴前进行探查,生成隐私风险报告。 |
| 5. 🧑🔬 阶段五:数据科学工作流中的特征探索 | 帮助数据科学家在建模前快速理解数据分布,进行特征工程。 | 快速EDA:快速获取目标数据集分布、缺失值和数据类型,加速模型探索-4。 整合工作流:可与Pandas DataFrame无缝集成,直接在Notebook中生成数据质量报告-4。 |
总的来说,DataProfiler 是用在 “先探查,后开发” 的思路里。
它会在项目前期帮你回答“我的数据里到底有什么?有什么风险?”这个问题,并在后期帮你持续监控“我的数据和上个月比,有没有发生奇怪的变化?”。因此,它通常在正式编写数据质量测试用例(如用Great Expectations)之前,以及在数据资产上线后持续运行的监控任务中发挥最大价值。
[toc]
这是一份关于 Flink CDC、Flume 和 DataX 这三个数据同步工具的对比分析,我会从核心定位、优势、适用场景以及项目实战中的常见问题这几个方面展开。
| 工具 | 核心定位 | 处理模式 | 典型场景 |
|---|---|---|---|
| Flink CDC | 实时数据捕获与同步 | 流处理 (Streaming) | 实时数据集成、构建实时数仓、数据库变更订阅-6 |
| Flume | 海量日志采集与传输 | 流式传输 (Streaming) | 日志收集、监控数据聚合、系统间事件中转- |
| DataX | 离线/批量数据同步 | 批处理 (Batch) | 异构数据源离线迁移、数据备份、T+1数据仓库构建– |
Flink CDC(Change Data Capture)的优势在于其全增量一体化的流处理架构。它基于Flink引擎,不仅能实时捕获数据库的增量变更日志(如MySQL的binlog),还能一次性完成历史数据的全量快照读取,并在全量同步完成后无缝切换到增量模式,整个过程无需中断,保证数据一致性-。这使其成为构建实时数据湖、实时数仓和需要低延迟响应的数据管道的理想选择-1。
Flume 的核心优势在于其专为日志而生的轻量级、高可靠架构。它采用简单的管道(Source-Channel-Sink)模型,稳定且易于部署,是处理非结构化或半结构化日志数据的利器-10。因此,它广泛用于将分散在大量服务器上的应用日志、系统日志、访问日志等汇聚到HDFS、Kafka等中央存储系统中,为后续的离线或实时分析提供原始数据–。
DataX 的优势在于其高度的可扩展性和对离线批量同步场景的专注。它采用插件化架构,官方支持大量异构数据源(关系型数据库、NoSQL、大数据存储等),通过配置JSON文件即可完成开发,简单稳定-6-20。因此,当需要进行大规模、周期性、对实时性无要求的离线数据迁移、数据备份或T+1数据仓库ETL时,DataX是一个非常成熟可靠的选择--20。
这三个工具在落地时也会遇到一些典型问题,以下是常见问题速查表,方便你快速定位和排查。
| 问题维度 | Flink CDC | Flume | DataX |
|---|---|---|---|
| 数据重复/丢失 | 高频问题。主要发生在全增量切换时,因位点(offset)管理或chunk边界重叠导致数据重复-47;或因Checkpoint未成功、重启时位点提交不及时导致-47。 | 可能因Source或Sink的事务配置不当导致数据重复或丢失,但日志场景通常对数据质量要求相对宽松。 | 任务中断后重启可能重跑部分数据,导致重复。不支持断点续传。 |
| 性能瓶颈 | 数据库连接压力:全量读取和实时CDC会与数据库建立长连接,增加源库负载-。 状态管理:处理大规模数据时对Flink状态后端(如RocksDB)要求高-。 | 通道吞吐量:单个Source-Sink对的处理能力有限,处理海量数据时容易成为瓶颈。 缓存开销:Channel的内存或文件缓冲可能带来额外延迟-11。 | 单机模式瓶颈:受限于单节点资源,处理TB级以上数据时性能可能不稳定-27。 资源竞争:并发任务多时,对CPU和内存的争用明显。 |
| 部署与运维 | 技术栈复杂:依赖Flink集群,需熟悉Flink作业的提交、调优和监控-6。 Schema变更:需处理源表结构变化,配置较繁琐。 | 配置繁琐:每个数据流都需独立配置,维护大量Flume Agent较复杂-11。 故障转移:Agent故障时的自动恢复和负载均衡需精心设计。 | JSON配置复杂:复杂的同步任务需要编写冗长的JSON文件,易出错,缺乏可视化界面-20-27。 依赖外部调度:本身无调度能力,需集成DolphinScheduler等工具。 |
| 数据一致性 | Exactly-Once保证:依赖Flink的Checkpoint机制,理论上可保证端到端的一致性,但在实际复杂环境中达成代价较高-。 | At-least-Once保证:默认设计倾向于至少一次,确保数据不丢失,但可能重复。 | 批量一致:以批为单位,一批数据要么全成功要么全失败,无事务性支持,易出现数据不一致。 |
| 环境依赖 | 强依赖Flink:必须有一套稳定、资源充足的Flink集群,且对Flink版本有要求。 | 依赖Hadoop生态:常用于与HDFS、HBase等集成,对Java环境有要求。 | 轻量级:只需Java环境,单机即可运行,对集群依赖小。 |
[toc]
表格
| 对比维度 | Spark | Flink |
|---|---|---|
| 计算模型 | 微批、离线为主 | 纯流式、事件驱动 |
| 延迟 | 分钟级、小时级 | 毫秒 / 秒级低延迟 |
| 时间语义 | 处理时间为主,事件时间弱 | 事件时间完善,水位线、乱序处理强 |
| 一致性 | 勉强支持 Exactly-Once | 原生支持 Exactly-Once |
| 窗口能力 | 简单窗口,迟到处理弱 | 滚动 / 滑动 / 会话 / 迟到数据全套支持 |
| 吞吐能力 | 超高吞吐,适合海量批量 | 吞吐不错,略低于 Spark 批处理 |
| 容错机制 | 血统依赖、重算整个分区 | Checkpoint + 状态增量快照,细粒度容错 |
| 状态管理 | 弱,依赖外部存储 | 原生状态后端(RocksDB)、TTL、增量 CK |
| 开发成本 | SQL 友好,上手简单 | API 稍复杂,实时场景功能更全 |
| 资源开销 | 批量调度资源利用率高 | 常驻任务,资源长期占用 |
| 小文件 | 容易产生,需手动合并 | 可窗口攒批 + 文件滚动策略控制 |
准实时(5~10 分钟延迟)
可用 Spark Streaming,但现在业界更倾向直接上 Flink,统一流批引擎。
流批一体架构
直接用 Flink 流批统一,一套代码跑实时 + 离线,不用维护两套引擎。
数据湖 Hudi/Iceberg
写入、增量消费 Flink 更适配;离线分析 Spark 更强。
技术选型主要看延迟要求、业务场景、数据量级、一致性诉求:
离线批量、海量 ETL、日级小时级数仓、高吞吐低成本场景选 Spark;
低延迟实时计算、实时大屏、实时特征、乱序日志、事件时间窗口、要求精确一次语义的场景选 Flink;
现在主流架构采用Flink 流批一体,实时离线一套引擎统一开发,离线分析依然用 Spark 做大规模批量查询和报表。
离线批量、数仓报表 → 选 Spark
实时低延迟、窗口乱序、精准一次 → 选 Flink
流批一体统一架构 → 全站 Flink,离线分析保留 Spark
答:Spark 是微批计算,把流拆成小批量离线任务跑;
Flink 是原生流式计算,一条一条事件持续处理,是真正的流引擎。
答:Spark 分钟级、小时级,适合离线;
Flink 毫秒 / 秒级低延迟,适合实时大屏、实时指标。
答:Flink 强很多,原生支持 EventTime 事件时间、Watermark 水位线、迟到数据处理;
Spark 偏向处理时间,事件时间、乱序支持较弱。
答:Flink 窗口完备:滚动、滑动、会话、全局窗口,支持水位线 + 允许迟到 + 侧输出兜底;
Spark 窗口简单,对乱序、晚到数据处理能力差。
答:Flink 原生 Checkpoint + 状态后端,天然支持 Exactly-Once;
Spark 只能靠业务幂等、事务兜底,引擎层面支持不彻底。
答:Flink 内置 RocksDB 状态后端、TTL 过期、增量 Checkpoint,适合长期驻留任务;
Spark 无原生状态,状态只能存在内存或外部存储,不适合复杂实时聚合。
答:Spark 基于 RDD 宽窄依赖血统重算,失败重算整个分区;
Flink 基于 Checkpoint 快照,只恢复失败节点最近状态,粒度更细、开销更小。
答:Spark 批量调度,吞吐更高、资源利用率好、成本低,适合海量离线;
Flink 任务常驻集群,资源长期占用,吞吐略低于 Spark 批处理。
答:Flink 靠水位线推进时间,三层处理:水位线等待 + 允许迟到 + 侧输出兜底;
Spark 对乱序不友好,很难保证时间窗口精度。
答:
离线数仓、批量 ETL、历史回溯、日 / 小时级报表、海量高吞吐低成本 → 选 Spark;
实时大屏、实时特征、实时风控、乱序日志、要求低延迟 + 精确一次 → 选 Flink;
流批一体架构统一用 Flink 做实时 + 离线,复杂离线分析保留 Spark。
表格
| 框架 | 计算模型 | 延迟 | 时间语义 | 窗口 / 乱序 | 状态管理 | 生产常用 |
|---|---|---|---|---|---|---|
| Spark Streaming | 固定间隔微批 | 秒~分钟级 | 处理时间 | 弱、不支持乱序 | 无原生状态 | 已淘汰 |
| Spark Structured Streaming | 默认微批;Continuous 真流 | 百 ms~ 秒级 | 一般支持事件时间 | 一般,水位线偏弱 | 有状态但偏弱 | 准实时、简单实时 |
| Flink | 原生流式事件驱动 | 毫秒~秒级 | 完善 EventTime | 水位线 + 迟到 + 侧输出全套 | 原生 RocksDB+TTL + 增量 CK | 高可用复杂实时首选 |
Flink (毫秒级 10ms) < Structured Streaming (百毫秒 100ms) < Spark Streaming (秒级 1000ms+)
Spark Streaming 是老旧固定微批已淘汰;
Structured Streaming 虽然 API 是流式,默认底层依旧微批轮询触发,延迟和乱序处理能力有限;
Flink 是原生事件驱动流式,支持完善的事件时间、水位线、迟到数据和原生状态管理,低延迟、一致性更好,复杂实时业务优先选 Flink。
简单准实时、熟悉 Spark 生态 → 用 Structured Streaming
低延迟、乱序严重、窗口精准、长期状态 → 必选 Flink
老项目 Spark Streaming 直接废弃重构
Spark Structured Streaming 默认还是微批,但从架构设计上伪装成了 “流式”,不是真正的事件逐条处理;只有开启 Continuous 模式才是真正低延迟流。
Spark Structured Streaming 底层依旧是按固定间隔生成 Batch,
每隔一段时间触发一次 Job,批量处理这段时间的数据。
有,但极少用:
Continuous 连续处理模式
1 | spark.sql.streaming.continuous.execution.enabled=true |
Spark Structured Streaming 默认底层仍是微批处理,按照固定时间间隔生成批次任务执行,并不是像 Flink 那样事件驱动的原生流式;
虽然 API 设计成流式编程模型,但调度引擎还是微批。只有开启 Continuous 连续模式才是真正低延迟流,但生产稳定性不足,一般不用。
[toc]
核心定位:本宝典适配机器学习、数据挖掘工程师全岗位(初级/中级/高级),聚焦面试高频考点、核心技能、项目话术与实操要点,兼顾理论深度与工程落地,帮你快速梳理核心知识、规避面试误区,高效通关各类企业面试(互联网、AI公司、传统行业数字化岗均适用)。
核心原则:面试核心考察「理论基础+工程能力+项目落地+业务理解」,避免只背公式不聊落地,拒绝只讲概念不懂实操,突出“能建模、能落地、能解决业务问题”的核心竞争力。
面试官您好,我是XX,有X年机器学习/数据挖掘相关经验,熟练掌握机器学习算法(LR、树模型、深度学习等)、数据处理工具(Python、Spark、Pandas),主导/参与过XX(项目类型,如用户流失预测、推荐系统、异常检测)项目,擅长从业务场景出发,完成数据预处理、特征工程、模型训练、部署落地全流程。我的核心优势是XX(如:擅长特征工程、熟悉分布式建模、能快速定位模型问题),希望深耕XX(方向,如推荐、风控、NLP)领域,助力业务价值落地。
面试官您好,我是XX,毕业于XX院校XX专业,有X年机器学习/数据挖掘工程师工作经验,主要聚焦XX(核心方向,如风控建模、用户画像、时序预测)领域。
技术层面,我熟练掌握:① 算法:逻辑回归、随机森林、XGBoost/LightGBM、CNN/RNN/Transformer等,懂算法原理、调参技巧及适用场景;② 工程工具:Python(Pandas、Scikit-learn、PyTorch/TensorFlow)、Spark(MLlib)、Hive,能处理千万级以上数据,完成分布式特征工程与模型训练;③ 落地能力:熟悉模型部署(FastAPI、TensorFlow Serving)、监控(数据漂移、模型精度)与迭代优化。
项目层面,我主导过XX项目(核心项目),从业务需求拆解、数据采集清洗,到特征工程、模型选型调优,再到上线部署,全程负责,最终实现XX(业务价值,如准确率提升15%、流失率降低10%)。
我注重“算法落地而非单纯调参”,擅长结合业务场景选择合适的模型,解决实际业务痛点,希望能加入贵司,在XX领域持续深耕,创造业务价值。
无论初级/中级/高级,面试均围绕以下4个维度考察,优先级:工程落地 > 算法基础 > 项目经验 > 业务理解,不同职级侧重不同:
面试聊项目,重点不是“调了多少参、用了什么算法”,而是“如何从业务出发,用算法解决实际问题,带来什么业务价值”,核心逻辑:业务痛点 → 技术方案 → 落地过程 → 效果复盘 → 优化方向。
项目名称:XX(如:基于LightGBM的用户流失预测系统)
结尾提示:面试的核心是“展现你的能力匹配岗位需求”,重点突出“能建模、能落地、懂业务”,结合本宝典的知识点和话术,针对性准备,即可高效通关。祝各位面试顺利!
[toc]
数据一致性的核心定义是:数据在全生命周期中,始终保持准确、完整、统一、符合业务规则的状态,在不同环节、不同系统、不同时间维度下,不存在矛盾、冲突、偏差、丢失或重复的情况。
它是数据仓库、数据库、分布式系统的核心基础指标,在不同场景下有不同的侧重点,其中数仓场景下的业务一致性,是你面试中需要重点掌握的核心内容。
这是一致性最原始的定义,也是所有数据系统的底层基础,对应数据库 ACID 四大特性中的Consistency(一致性)。
定义:事务执行前后,数据库的完整性约束、业务规则不会被破坏。
核心逻辑:一致性是事务的最终目的,原子性、隔离性、持久性都是为了保障一致性而存在的。
通俗例子:
银行转账场景,A 账户给 B 账户转 100 元,事务执行后,A 账户扣 100 元、B 账户加 100 元,两个账户的总金额必须和转账前完全一致,不能出现 A 扣了钱、B 没收到,或者总金额对不上的情况,这就是事务层面的一致性。
这是你之前数仓面试宝典中反复提到的核心概念,也是企业数仓建设的核心痛点,本质是保证数据在数仓全链路中,业务含义、统计口径、数值结果始终统一,不出现 “报表打架、指标对不上” 的问题。
主要分为 5 个核心维度,覆盖离线 + 实时全场景:
[toc]
面试原题:百亿条有序时间数据,怎么最快查到某条记录?
原理背诵:有序数组,每次取中间值比较,折半缩小范围;时间复杂度O(logN)。
大数据场景:Hive分区查找、索引查询、Doris有序列快速过滤、日志时间检索。
面试话术:海量有序数据优先二分查找,比遍历快得多;数仓中用于分区裁剪、有序索引定位。
原理:快慢指针、左右指针;一次遍历完成去重、合并、区间筛选。
大数据场景:有序大文件合并、日志去重、连续时间段筛选、行为轨迹分析。
面试举例:两个有序超大日志文件,双指针一趟遍历合并为一个有序文件。
原理背诵:单向链表、双向链表;插入删除快、查询慢;LRU底层=哈希表+双向链表。
大数据场景:Redis缓存、Flink状态管理、Kafka偏移量链表维护。
栈:后进先出;用于括号匹配、递归回溯、任务回滚。
队列:先进先出;Kafka消息队列、任务排队、流量削峰。
阻塞队列:Flink线程池、Spark资源调度、生产消费模型。
原理:自己调用自己,拆分重复子问题;注意递归深度防止栈溢出。
数仓场景:商品类目层级、地区树形结构、部门层级、血缘溯源。
面试必考:前序、中序、后序、层序遍历。
大数据场景:数据血缘树、任务依赖树、菜单类目树、索引B+树。
面试话术:数据平台血缘分析底层采用二叉树遍历,递归追溯上游依赖表。
面试原题:100个有序大文件,合并成一个全局有序文件?
原理:每个文件保留一个指针,最小堆维护最小值;每次取出最小写入结果。
生产场景:Hive合并小文件、日志合并、离线大排序、历史数据规整。
原理:Hash统计频次,结合TopK找出高频key。
数仓用途:倾斜key排查、热点用户、爆款商品、恶意IP限流。
结构背诵:时间戳+机器码+序列号;全局唯一、趋势递增、无重复。
大数据场景:订单ID、埋点日志ID、分布式数据表主键、Flink唯一标识。
场景:Kafka分区分发、Redis集群、服务网关、集群负载。
生产场景:埋点日志限流、风控防刷、Kafka流量削峰。
原理:字符逐层存储,公共前缀共享节点;查询极快、节省内存。
场景:URL黑名单、敏感词过滤、日志前缀筛选、域名匹配。
数仓用途:清洗手机号、身份证、URL、特殊乱码、脏数据过滤;DWD层高频使用。
面试话术:哈希算法将任意长度数据转为定长哈希值;用于去重、分片、加密、路由。
生产应用:shuffle分区、分桶、一致性哈希、布隆过滤器、文件校验。
面试原题:Kafka分区为什么能保证同一个key发送到同一个分区?扩容为什么不会大量迁移数据?
原理背诵:
大数据应用:Kafka分区路由、Redis集群、分库分表、数据分片。
适用场景:用户随机抽样、风控随机打散。
算法思想:从后往前遍历,当前位置与随机位置交换,时间复杂度O(n)。
面试背诵话术:Fisher-Yates洗牌,保证概率均匀、无偏,大数据抽样常用,Hive中order by rand()底层就是洗牌算法。
面试原题:亿级数据随机抽取100条,内存放不下全部数据怎么做?
原理背诵:
使用场景:Hive抽样、日志随机采样、风控黑名单抽样。
原理话术:数据倾斜本质是hash分区不均匀,热点key落在同一个reduce。解决方案:空值打散、加盐随机、热点key拆分、局部聚合。
面试原题:Redis缓存满了怎么淘汰?
原理背诵:最近最少使用淘汰,链表维护访问顺序,头部最新、尾部最久未使用;淘汰尾部节点;Redis近似LRU。
数仓场景:Redis缓存维度表、用户画像、热点特征,使用LRU淘汰冷数据。
面试话术:二进制数组+多个哈希函数;判断一定不存在、不一定存在;优点内存极小;缺点无法删除、存在误判。
业务场景:用户去重、黑名单过滤、日志重复判断。
面试原题:10亿条访问日志,找出访问量最高的100个IP,内存放不下全部数据怎么做?
算法原理(背诵):
大数据生产应用:热门商品、热点IP、倾斜key排查、流量TOP排行。
面试话术:海量数据求TopK不能全局排序,时间复杂度太高;使用小顶堆,时间复杂度O(nlogK),内存只保留K个元素,适合超大数据量。
面试原题:1亿用户ID,判断用户是否登录,怎么极致节省内存?
原理背诵:
数仓应用场景:日活用户去重、签到状态、用户黑白名单、用户留存标记。
面试原题:Flink滚动窗口、滑动窗口底层原理?
算法原理:
业务场景:实时最近1小时交易额、最近30秒流量、实时告警、滑动留存。
原理背诵:选取基准值,小于基准放左边、大于基准放右边;递归拆分;平均时间复杂度O(nlogn)。
大数据优化点:Hive、Spark底层排序都是优化快排;结合外排序解决磁盘海量数据排序。
面试话术:大数据生产不用冒泡,全部使用快速排序,无序海量数据效率最高。
面试原题:一个100G日志文件,内存只有4G,怎么排序?
解题步骤(必背):
生产应用:Hive大表排序、离线超大日志排序、历史数据规整。
面试原题:十亿条交易金额,求交易中位数,内存放不下?
解法话术:双堆法,大顶堆存前半段、小顶堆存后半段;堆顶即为中位数;海量数据配合分片抽样,近似中位数。
原理:每一步只做当前最优选择,不回溯;局部最优推全局最优。
大数据场景:Azkaban调度资源分配、任务优先级、压缩策略、冷热数据存储选型。
原理背诵:
生产作用:降低存储成本、提升查询速度、优化集群压力。
1、海量数据处理:蓄水池、位图、布隆、外排序、多路归并、TopK
2、分片路由算法:一致性哈希、哈希取模、负载均衡
3、流式实时算法:滑动窗口、水位线、限流、令牌桶
4、基础高频算法:快排、二分、双指针、递归、栈队列、链表
5、工程实用算法:雪花ID、字典树、正则、贪心、冷热分离
6、去重算法三件套:BitMap > 布隆过滤器 > Hash表
说明:全部最短口诀,不记原理、不记代码,面试张口就说,专门针对大数据数仓开发。
海量抽样蓄水池,去重位图布隆池;
分片一致哈希环,实时窗口水位齐;
堆求top快排序,二分指针最简单;
倾斜加盐打散用,冷热归档省机器;
工程雪花唯一id,缓存LRU永不弃。
1 | row_number() over(partition by 分组字段 order by 排序字段 desc) rn |
1 | select * from ( |
1 | sum(col) over(order by dt rows between unbounded preceding and current row) |
1 | count(1)/sum(count(1)) over(partition by class_id) |
业务场景:每个部门薪资最高前2人、每个商品类目销量TOP3。
标准答案SQL:
1 | select * from ( |
1 | SELECT |
1 | AVG(price) OVER(ORDER BY dt ROWS BETWEEN 2 PRECEDING AND CURRENT ROW) |
1 | select distinct user_id from log; |
1 | select user_id,max(dt) from log group by user_id; |
1 | select * from ( |
面试话术:生产环境禁止大表distinct,容易触发数据倾斜,优先row_number分组去重。
Hive/Spark SQL
1 | SELECT |
面试话术背诵:行转列使用collect_list无序聚合、collect_set去重聚合,搭配concat_ws拼接字符串,常用于标签合并、多属性合并。
1 | SELECT |
面试话术背诵:使用lateral view炸裂函数,搭配explode数组拆分,split切割字符串,实现列转行,常用于标签拆分、多维拆解。
1 | -- 原写法倾斜 |
1 | -- 第一层:局部聚合加盐 |
/*+ BROADCAST(small_table) */1 | SELECT /*+ BROADCAST(b) */ a.* |
思路:日期减去行号,相同即为连续
1 | # 连续登陆超过3天 |
需求:计算每日新增用户、次日留存、7日留存。
留存定义:当天新增用户,后续某天再次活跃。
面试话术背诵:留存使用自关联,当天新增表关联未来活跃表;离线留存T+1计算,实时留存使用Flink状态做当日留存。
标准答案SQL:
1 | # 次日留存 |
业务:浏览-加购-下单-支付,每一步转化率。
解题思想:同一用户行为路径、判断是否走完下一个节点。
1 | select |
面试话术:漏斗核心逻辑是用户行为埋点、行为编号,分层统计人数,计算转化率;大厂一般使用Flink实时漏斗、离线Hive漏斗。
开窗函数 rows between 边界,生产高频。
1 | select |
关键字背诵:unbounded preceding(首行)、current row(当前行)。