0%

1 Ray 分布式计算核心概念

一句话定位:Ray 是面向 AI/ML 的分布式框架,用 Task/Actor 统一抽象 + Plasma 共享内存 + 细粒度动态调度,在异构资源(CPU/GPU)上跑 AI 数据处理与推理。面试核心是”为什么 AI 数据链路常比 Spark 更合适”。

2 整体架构:节点与核心组件

Ray 集群由若干节点组成,每个节点都有一个 Raylet(本地调度器)和 Object Store(基于 Plasma 的共享内存对象存储),集群层面还有全局的 GCSDriver 是运行用户主程序的进程。

Ray 系统架构

  • GCS(Global Control Store):全局控制存储,保存系统级元数据——节点信息、对象所在位置、函数/类定义等;节点之间不直接互相通信,而是通过 GCS 共享状态,避免中心化调度器的单点瓶颈。
  • Raylet:每个节点一个的本地调度器,负责资源管理、任务排队、依赖对象拉取与 Worker 生命周期管理,是”自底向上调度”的执行者。
  • Worker 进程:真正执行 Task / 运行 Actor 的进程,按 CPU/GPU 资源需求由 Raylet 拉起,粒度可细到小数(如 0.3 GPU)。
  • Object Store(Plasma):基于共享内存的本地对象存储,节点内零拷贝、节点间按需拉取。

3 Task(无状态)vs Actor(有状态)

  • Task@ray.remote 修饰的普通函数,无状态、输入→输出通过 ObjectRef 返回,调度器动态分配到空闲节点。适合无状态并行(map、批量推理)。
  • Actor@ray.remote 修饰的类,在固定节点上常驻一个进程、持有状态;其方法调用串行执行。适合有状态服务——参数服务器、模型推理服务、累积计数/状态机。
  • 关键区别:Task 调用间不共享内存、可任意调度;Actor 把状态”钉”在节点上,调用有顺序保证但需考虑状态热点。

Task vs Actor 对比

补充几点:

  • Task 每次调用可能落在不同节点/Worker,函数本身不保存任何跨调用状态,返回值通过 ObjectRef 异步取回。
  • Actor 创建后固定在某个节点,其方法(method)按提交顺序串行执行,天然避免并发写同一状态的竞态;若需并发,可用 AsyncActor(方法声明为 async def)。
  • 状态热点:所有调用都打到同一 Actor 进程,可能成为吞吐瓶颈,通常需要分片(多个 Actor 各持一部分状态)。

4 分布式对象存储 Plasma

  • Plasma 是基于共享内存的分布式对象存储:大对象(DataFrame、模型权重、中间 tensor)以零拷贝方式在节点内/节点间传递,避免反复序列化拷贝。
  • 对象由 ObjectRef 引用,生命周期由引用计数管理;这是 Ray 能在 Python 生态里高效搬运数据的基础。
  • 生命周期细节:每个对象由创建它的进程(owner)负责,采用分布式引用计数(distributed reference counting),引用归零才回收——既避免跨节点大对象被过早释放,也不会内存泄漏。
  • 跨节点语义:当任务依赖远端对象时,Raylet 会先把对象拉取(pull) 到本地 Object Store,之后本地 Worker 直接共享内存读取,避免反复走网络。

5 动态任务图调度

  • 与 Spark 的”提交时静态 DAG”不同,Ray 的任务在运行时动态生成,调度器(全局 GCS + 局部)按资源可用性把细粒度 task 放到合适节点(可精确到 0.1 GPU)。
  • 细粒度 → 更好的负载均衡与资源利用率,也天然支持不规则、依赖运行时的计算图。

自底向上(bottom-up)调度:Ray 没有全局中心调度器,而是由 Driver/Worker 把任务提交到本地 Raylet,本地 Raylet 先到 GCS 查输入对象在哪个节点,再决定在本地执行还是”转交”给持有数据的远端 Raylet。这种”数据/资源在哪里,就往哪里调”的方式既避免中心调度器瓶颈,又天然贴合数据本地性。

Ray 自底向上调度流程

6 容错:基于血缘(Lineage)的重执行

  • Ray 默认用血缘(lineage)记录任务依赖链(哪个对象由哪个 task 产生)。某个 task 或对象丢失时,调度器沿血缘重新执行上游 task 来重建,而不是像 Spark 那样依赖检查点快照。
  • 对象重建:Plasma 中的对象若因节点宕机丢失,可从血缘重新计算得到。
  • Actor 例外:Actor 是有状态长驻进程,无法简单”重执行”恢复,需要应用层通过 checkpoint 或重建 Actor 处理。这是 Ray 容错的薄弱点,也是面试常问的”Ray 容错边界”。
  • 生产中可借助高层库(Ray Train / Ray Serve)的 checkpoint 机制增强容错。

7 面试重点:为什么 AI 数据链路 Ray 常优于 Spark

  1. 异构资源:Ray 原生细粒度分配 GPU(如 0.3 GPU 给算子 A、0.7 给算子 B);Spark 以 executor/core 粗粒度调度,GPU 支持弱。
  2. Python 原生:Ray 围绕 NumPy/PyTorch 设计,借助 Plasma 零拷贝传 tensor;Spark 以 JVM 为主,PySpark UDF 经 py4j 序列化开销大。
  3. 动态/细粒度调度:AI 管线常含 Actor 服务、动态 batch、运行期才确定的依赖,Spark 静态 DAG 不灵活。
  4. 典型场景:离线批量推理、超参搜索、RL、预处理+训练一体。

参考:

  • 论文《Ray: A Distributed Framework for Emerging AI Applications》(Moritz et al., 2018, arXiv:1712.05889);
  • Anyscale《Offline Batch Inference: Comparing Ray, Apache Spark, and SageMaker》

Spark 内存管理与 GC 调优

一句话定位:Executor 内存怎么分、谁来借、GC 为什么抖,是 Spark 性能调优的地基;理解内存模型才能解释”为什么缓存一多就 Full GC”和”为什么分区太小反而更慢”。

1. Executor 内存模型

参考: Spark Executor内存管理

Executor 的 JVM 堆被划分为(以 spark.memory.fraction 默认 0.6 为例,扣除 300MB Reserved 之后):

  • 统一内存(Unified Memory):执行内存(Execution,排序/聚合/Join 的 shuffle 缓冲)与存储内存(Storage,缓存的 RDD/广播变量)共享一个池,默认各占 50%(spark.memory.storageFraction)。
    • 动态借用:执行内存不足时可向存储内存借空间(逐出缓存的 block);但存储内存不能强行逐出正在被使用的执行内存。即执行内存优先,这也是”缓存过多反而拖慢计算”的根因——存储占了统一池,执行只能被挤压或 spill 到磁盘。
  • 用户内存(User Memory):算子/UDF 自身对象、数据结构,不受 Spark 管理。
  • 堆外内存(Off-Heap)spark.memory.offHeap.enabled=true + spark.memory.offHeap.size;Tungsten 的二进制格式(排序/聚合/缓存)可直接写在堆外,由 OS 管理、不受 JVM GC 管辖,显著降低 GC 压力,但需手动设上限。

Spark JVM On-heap 内存分配
JVM On-heap 内存分配

2. GC 频繁的典型成因

  • 分区过小 → 单分区对象少但分区数爆炸 → Young GC 次数飙升;
  • 缓存过量(RDD/广播变量/大表缓存)→ Old Gen 持续高位 → Full GC;
  • 大批量/大对象(如超大 collect、宽行)→ 大对象直入 Old Gen,触发提前 Full GC;
  • 堆外未开 + 统一内存被存储挤占 → 执行内存 spill 到磁盘,反而增加对象生命周期。

3. 调优手段

  • 调分区数spark.sql.shuffle.partitions / repartition,让单任务数据量适中,减少对象数与 GC 频率。
  • 控制缓存:非必要不 cache;用 MEMORY_AND_DISK 而非纯内存;及时 unpersist
  • 开堆外:大数据量排序/聚合场景开启 Off-Heap,把 GC 压力移出 JVM。
  • GC 器与参数:G1GC(-XX:+UseG1GC)配合合理 -Xmn / -XX:InitiatingHeapOccupancyPercent;监控 GC 日志与 Spark UI Executors 页的 GC Time
  • 加资源:在单任务数据量确实大时,增大 executor.memory 或降 executor.cores(减少同进程并发对象)。

参考:Apache Spark 官方文档 Memory Management Overview / Tuning Guide

资源调度与执行计划优化

一句话定位:作业跑得快不快,一半在资源撮合(能不能拿到合适的容器),一半在执行计划(会不会自适应地选对分区数和 Join 策略)。Spark 3 的 AQE 把后者从”人肉调参”变成了”运行时自适应”。

1. YARN / K8s 资源撮合

  • YARN:ResourceManager 按 Container 申请(内存 + vcore)在 NodeManager 上分配;Capacity/Fair 调度器决定队列配额与共享;受 min-allocation / increment 规整化影响,实际占用常大于申请(见 Spark 资源评估)。
  • K8s:executor 以 Pod 申请,由 kube-scheduler 按节点资源/亲和性调度;支持动态资源分配与 Spot/抢占,弹性更好。

2. 任务优先级与抢占

  • 队列/作业设优先级,高优先级优先获得资源。
  • 抢占(Preemption):Fair/Capacity 模式下,在优雅终止期内回收低优先级容器的资源给高优先级,保障关键作业 SLA。
  • 代价:被抢占作业需重试/重算对应 partition,设计上要容忍。

3. AQE(自适应查询执行,Spark 3.0+)

运行时基于 shuffle 统计动态调整执行计划,三项核心能力:

  1. 分区合并(Coalesce Partitions):把多个小 post-shuffle 分区合并,避免小 task 过多、调度开销爆炸(无需手动设 spark.sql.shuffle.partitions)。
  2. 倾斜处理(Skew Join):检测到倾斜分区(> 中位数 N 倍且超阈值)自动拆分为多 task 并行(见 1.1.2 数据倾斜)。
  3. Join 策略切换:运行时若某表小于 spark.sql.autoBroadcastJoinThreshold,把 SortMergeJoin 动态转 BroadcastHashJoin,免 Shuffle。
  • 收益:对数据量/分布不确定的 ad-hoc 查询尤其明显;把”凭经验调参”变为”运行时自适应”。

4. 执行计划优化

  • explain() 读懂物理计划;开启 CBO(需列统计/直方图)让优化器选更优 Join 顺序;
  • 配合广播阈值、谓词下推、列裁剪,减少 shuffle 与扫描量。

参考:Apache Spark 官方文档 Adaptive Query Execution / Running Spark on YARN / Kubernetes

构建本地模拟环境

hudi代码走读

创建本地模拟环境,一步一步调试

主程序

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
String sourceTable = "source_01";
ParameterTool pt = ParameterTool.fromArgs(args);
String database = pt.get("database", "ods_glory");
String table = pt.get("table", "ods_glory_t1");
String hudiDemoHome = "apps/data/hudi";
String tablePath = String.format("%s/%s/%s.db/%s", System.getProperty("user.home"), hudiDemoHome, database, table);
log.info("userParameter: {}", pt.toMap());
//
StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
env.setParallelism(1);
env.enableCheckpointing(10000, CheckpointingMode.EXACTLY_ONCE);
Configuration configuration = tableEnv.getConfig().getConfiguration();
configuration.setString("table.dynamic-table-options.enabled", "true");
configuration.setString(PipelineOptions.NAME.key(), "flake-hudi-1.13");
log.info("envConfig: {}.", configuration.toMap());
//
tableEnv.executeSql(sourceDdl());
String targetDal = sinkTableDdl(table, tablePath, database);
log.info("targetDal:{}.", targetDal);
tableEnv.executeSql(targetDal);
// tableEnv.executeSql(String.format("select *,DATE_FORMAT(ts, 'yyyyMMdd') as dt from %s", table, sourceTable))
// System.out.println(tableEnv.explain(tableEnv.from(table)));
String sql = String.format("insert into %s select *,DATE_FORMAT(ts, 'yyyyMMdd') as dt from %s", table,
sourceTable);
tableEnv.executeSql(sql);
System.out.println(env.getExecutionPlan());

创建数据源

1
2
3
4
5
6
7
8
9
10
11
12
13
14
public static String sourceDdl() {
String sourceTable = "source_01";
return String.format("CREATE TABLE %s (\n" +
" uuid varchar(20),\n" +
" facc_time STRING ,\n" +
" facc_time_rongduan STRING,\n" +
" facc_type BIGINT ,\n" +
" fact_info STRING ,\n" +
" ts timestamp(3)\n" +
") WITH (\n" +
" 'connector' = 'datagen',\n" +
" 'rows-per-second' = '100'\n" +
")", sourceTable);
}

创建hudi表

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
public static String sinkTableDdl(String targetTable, String basePath, String dbName) {
return String.format("create table %s(\n" +
" uuid STRING,\n" +
" facc_time STRING ,\n" +
" facc_time_rongduan STRING,\n" +
" facc_type BIGINT ,\n" +
" fact_info STRING ,\n" +
" ts timestamp(3),\n" +
" dt string,\n" +
" PRIMARY KEY(uuid) NOT ENfORCED" +
")\n" +
" PARTITIONED BY (`dt`)\n" +
" with (\n" +
" 'connector' = 'hudi',\n" +
" 'path' = '%s', -- 替换成的绝对路径\n" +
" 'table.type' = 'MERGE_ON_READ',\n" +
" 'write.bucket_assign.tasks' = '8',\n" +
" 'write.tasks' = '8',\n" +
" 'write.operation' = 'upsert', -- upsert/insert\n" +
" 'changelog.enabled' = 'true',\n" +
" 'read.streaming.enabled' = 'true',\n" +
" 'read.streaming.check-interval' = '1',\n" +
" 'compaction.tasks' = '8',\n" +
" 'compaction.trigger.strategy'='num_commits',\n" +
" 'compaction.delta_commits' ='5',\n" +
" 'compaction.max_memory' = '3096',\n" +
" 'clean.retain_commits' = '30',\n" +
" 'hive_sync.enable' = 'false',\n" +
" 'hive_sync.mode' = 'hms',\n" +
" 'hive_sync.db' = '%s',\n" +
" 'hive_sync.table' = '%s',\n" +
" 'hive_sync.metastore.uris' = '%s'\n" +
")", targetTable, basePath, dbName, targetTable, "metastoreUrl");
}

go on,跟进调试

// TODO 整理到 flink demo

img

img

org.apache.hudi.table.HoodieTableSink

1
2
3
4
5
6
7
8
9
10
11
// bootstrap
final DataStream<HoodieRecord> hoodieRecordDataStream =
Pipelines.bootstrap(conf, rowType, parallelism, dataStream, context.isBounded(), overwrite);
// write pipeline
DataStream<Object> pipeline = Pipelines.hoodieStreamWrite(conf, parallelism, hoodieRecordDataStream);
// compaction
if (StreamerUtil.needsAsyncCompaction(conf)) {
return Pipelines.compact(conf, pipeline);
} else {
return Pipelines.clean(conf, pipeline);
}

img

org.apache.flink.table.planner.plan.nodes.exec.common.CommonExecSink#createSinkTransformation

img

经历三步: bootstrap → hoodieStreamWrite → compact / clean

bootstrap

  1. row_data_to_hoodie_record过程如下:

Flink rowData 转换为 HoodieRecord: org.apache.hudi.sink.transform.RowDataToHoodieFunction#toHoodieRecord

数据最终类型为: HoodieAvroRecord

  1. 如果配置了 index.bootstrap.enabled,会增加一个 index_bootstrap 节点,用于在flink state中保存instant信息(lastInstantTime)

StreamWrite

org.apache.hudi.sink.utils.Pipelines#hoodieStreamWrite

compact

clean

org.apache.hudi.sink.utils.Pipelines#bootstrap(org.apache.flink.configuration.Configuration, org.apache.flink.table.types.logical.RowType, int, org.apache.flink.streaming.api.datastream.DataStream<org.apache.flink.table.data.RowData>, boolean, boolean)

img

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
org.apache.hudi.sink.transform.RowDataToHoodieFunction

org.apache.hudi.sink.bootstrap.BootstrapOperator 这个需要继续分析

// org.apache.hudi.sink.utils.Pipelines#hoodieStreamWrite
dataStream
.keyBy(HoodieRecord::getRecordKey) // hoodieKey 需要单独分析
.transform("bucket_assigner", TypeInformation.of(HoodieRecord.class), new KeyedProcessOperator(new BucketAssignFunction(conf)))
// org.apache.hudi.sink.partitioner.BucketAssignFunction#processRecord
// hoodieKey 的 recordKey 和 partitionPath 代表是什么意思?
// org.apache.hudi.common.model.HoodieRecordLocation 和 org.apache.hudi.common.model.HoodieRecordGlobalLocation
.uid("uid_bucket_assigner_" + conf.getString(FlinkOptions.TABLE_NAME))
.setParallelism((Integer)conf.getOptional(FlinkOptions.BUCKET_ASSIGN_TASKS).orElse(defaultParallelism))
.keyBy((record) -> { return record.getCurrentLocation().getFileId(); })
.transform("stream_write", TypeInformation.of(Object.class), operatorFactory)
// org.apache.flink.streaming.api.operators.SimpleUdfStreamOperatorFactory
// Flink的单UDF的StreamOperatorFactory
// org.apache.hudi.sink.common.WriteOperatorFactory#createStreamOperator
.uid("uid_stream_write" + conf.getString(FlinkOptions.TABLE_NAME))
.setParallelism(conf.getInteger(FlinkOptions.WRITE_TASKS));

dataStream
.transform("compact_plan_generate", TypeInformation.of(CompactionPlanEvent.class), new CompactionPlanOperator(conf))
.setParallelism(1)
.rebalance()
.transform("compact_task", TypeInformation.of(CompactionCommitEvent.class), new ProcessOperator(new CompactFunction(conf)))
.setParallelism(conf.getInteger(FlinkOptions.COMPACTION_TASKS))
.addSink(new CompactionCommitSink(conf))
.name("compact_commit")
.setParallelism(1);
1
2
3
4
5
// org.apache.hudi.sink.partitioner.BucketAssignFunction#processRecord
HoodieRecord<?> deleteRecord = new HoodieAvroRecord(new HoodieKey(recordKey, oldLoc.getPartitionPath()), this.payloadCreation.createDeletePayload((BaseAvroPayload)record.getData()));
deleteRecord.setCurrentLocation(oldLoc.toLocal("U"));
deleteRecord.seal();
out.collect(deleteRecord);

// stream_write
// org.apache.hudi.sink.StreamWriteFunction
org.apache.hudi.sink.StreamWriteFunction#bufferRecord

FLIP-188

FLIP-188翻译

聊聊 Flink Table Store:流批一体存储最新进展|InfoQ大会早班车第16期

Flink Table Store v0.2 应用场景和核心功能

Paimon 文件结构

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
tree -h
[ 192] .
├── [1.3K] bucket-0
│   ├── [6.6M] data-271be4ad-8396-4aa5-9b1d-1de73a86c279-0.orc
│   ├── [8.8M] data-271be4ad-8396-4aa5-9b1d-1de73a86c279-1.orc
│   ├── [8.9M] data-271be4ad-8396-4aa5-9b1d-1de73a86c279-10.orc
│   ├── [8.8M] data-271be4ad-8396-4aa5-9b1d-1de73a86c279-11.orc
│   ├── [8.9M] data-271be4ad-8396-4aa5-9b1d-1de73a86c279-12.orc
│   ├── [8.8M] data-271be4ad-8396-4aa5-9b1d-1de73a86c279-2.orc
│   ├── [8.8M] data-271be4ad-8396-4aa5-9b1d-1de73a86c279-3.orc
│   ├── [8.8M] data-271be4ad-8396-4aa5-9b1d-1de73a86c279-4.orc
│   ├── [8.8M] data-271be4ad-8396-4aa5-9b1d-1de73a86c279-5.orc
│   ├── [8.9M] data-271be4ad-8396-4aa5-9b1d-1de73a86c279-6.orc
│   ├── [8.8M] data-271be4ad-8396-4aa5-9b1d-1de73a86c279-7.orc
│   ├── [8.9M] data-271be4ad-8396-4aa5-9b1d-1de73a86c279-8.orc
│   ├── [8.8M] data-271be4ad-8396-4aa5-9b1d-1de73a86c279-9.orc
│   ├── [ 92M] data-6608c270-fec6-4dcc-903d-f752501191af-0.orc
│   ├── [ 84M] data-6608c270-fec6-4dcc-903d-f752501191af-1.orc
│   ├── [128M] data-6608c270-fec6-4dcc-903d-f752501191af-2.orc
│   ├── [115M] data-6608c270-fec6-4dcc-903d-f752501191af-3.orc
│   ├── [ 84M] data-6608c270-fec6-4dcc-903d-f752501191af-4.orc
│   ├── [ 19M] data-6608c270-fec6-4dcc-903d-f752501191af-5.orc
│   ├── [ 28M] data-c991b192-ff5e-4654-a62f-18ad51821d0a-0.orc
│   ├── [ 29M] data-c991b192-ff5e-4654-a62f-18ad51821d0a-1.orc
│   ├── [ 29M] data-c991b192-ff5e-4654-a62f-18ad51821d0a-10.orc
│   ├── [550K] data-c991b192-ff5e-4654-a62f-18ad51821d0a-11.orc
│   ├── [ 29M] data-c991b192-ff5e-4654-a62f-18ad51821d0a-12.orc
│   ├── [ 29M] data-c991b192-ff5e-4654-a62f-18ad51821d0a-13.orc
│   ├── [ 29M] data-c991b192-ff5e-4654-a62f-18ad51821d0a-14.orc
│   ├── [ 29M] data-c991b192-ff5e-4654-a62f-18ad51821d0a-15.orc
│   ├── [666K] data-c991b192-ff5e-4654-a62f-18ad51821d0a-16.orc
│   ├── [ 3] data-c991b192-ff5e-4654-a62f-18ad51821d0a-17.orc
│   ├── [614K] data-c991b192-ff5e-4654-a62f-18ad51821d0a-2.orc
│   ├── [ 29M] data-c991b192-ff5e-4654-a62f-18ad51821d0a-3.orc
│   ├── [ 29M] data-c991b192-ff5e-4654-a62f-18ad51821d0a-4.orc
│   ├── [ 29M] data-c991b192-ff5e-4654-a62f-18ad51821d0a-5.orc
│   ├── [ 29M] data-c991b192-ff5e-4654-a62f-18ad51821d0a-6.orc
│   ├── [ 22M] data-c991b192-ff5e-4654-a62f-18ad51821d0a-7.orc
│   ├── [ 29M] data-c991b192-ff5e-4654-a62f-18ad51821d0a-8.orc
│   ├── [ 29M] data-c991b192-ff5e-4654-a62f-18ad51821d0a-9.orc
│   ├── [ 40M] data-f9dea92c-9e26-4966-ab62-3ed42a5c41c2-0.orc
│   ├── [ 34M] data-f9dea92c-9e26-4966-ab62-3ed42a5c41c2-1.orc
│   └── [100M] data-f9dea92c-9e26-4966-ab62-3ed42a5c41c2-2.orc
├── [2.2K] manifest
│   ├── [1.8K] manifest-a5fcc42f-c458-4a3e-a0c3-ce8780289413-0
│   ├── [1.8K] manifest-a5fcc42f-c458-4a3e-a0c3-ce8780289413-1
│   ├── [1.8K] manifest-a5fcc42f-c458-4a3e-a0c3-ce8780289413-10
│   ├── [2.2K] manifest-a5fcc42f-c458-4a3e-a0c3-ce8780289413-11
│   ├── [1.8K] manifest-a5fcc42f-c458-4a3e-a0c3-ce8780289413-12
│   ├── [1.8K] manifest-a5fcc42f-c458-4a3e-a0c3-ce8780289413-13
│   ├── [1.8K] manifest-a5fcc42f-c458-4a3e-a0c3-ce8780289413-14
│   ├── [1.8K] manifest-a5fcc42f-c458-4a3e-a0c3-ce8780289413-2
│   ├── [1.8K] manifest-a5fcc42f-c458-4a3e-a0c3-ce8780289413-3
│   ├── [1.8K] manifest-a5fcc42f-c458-4a3e-a0c3-ce8780289413-4
│   ├── [1.8K] manifest-a5fcc42f-c458-4a3e-a0c3-ce8780289413-5
│   ├── [2.2K] manifest-a5fcc42f-c458-4a3e-a0c3-ce8780289413-6
│   ├── [1.8K] manifest-a5fcc42f-c458-4a3e-a0c3-ce8780289413-7
│   ├── [1.8K] manifest-a5fcc42f-c458-4a3e-a0c3-ce8780289413-8
│   ├── [1.8K] manifest-a5fcc42f-c458-4a3e-a0c3-ce8780289413-9
│   ├── [2.0K] manifest-fe1ff04c-e66a-4532-88c7-0f2c34d261d8-0
│   ├── [2.5K] manifest-fe1ff04c-e66a-4532-88c7-0f2c34d261d8-1
│   ├── [2.2K] manifest-fe1ff04c-e66a-4532-88c7-0f2c34d261d8-2
│   ├── [2.2K] manifest-fe1ff04c-e66a-4532-88c7-0f2c34d261d8-3
│   ├── [2.1K] manifest-fe1ff04c-e66a-4532-88c7-0f2c34d261d8-4
│   ├── [2.5K] manifest-fe1ff04c-e66a-4532-88c7-0f2c34d261d8-5
│   ├── [2.2K] manifest-fe1ff04c-e66a-4532-88c7-0f2c34d261d8-6
│   ├── [2.2K] manifest-fe1ff04c-e66a-4532-88c7-0f2c34d261d8-7
│   ├── [ 910] manifest-list-4ad2bd0d-79be-4012-b55f-bdc916deac2d-0
│   ├── [ 769] manifest-list-4ad2bd0d-79be-4012-b55f-bdc916deac2d-1
│   ├── [1000] manifest-list-4ad2bd0d-79be-4012-b55f-bdc916deac2d-10
│   ├── [ 769] manifest-list-4ad2bd0d-79be-4012-b55f-bdc916deac2d-11
│   ├── [1012] manifest-list-4ad2bd0d-79be-4012-b55f-bdc916deac2d-12
│   ├── [ 769] manifest-list-4ad2bd0d-79be-4012-b55f-bdc916deac2d-13
│   ├── [1021] manifest-list-4ad2bd0d-79be-4012-b55f-bdc916deac2d-14
│   ├── [ 769] manifest-list-4ad2bd0d-79be-4012-b55f-bdc916deac2d-15
│   ├── [ 957] manifest-list-4ad2bd0d-79be-4012-b55f-bdc916deac2d-2
│   ├── [ 769] manifest-list-4ad2bd0d-79be-4012-b55f-bdc916deac2d-3
│   ├── [ 969] manifest-list-4ad2bd0d-79be-4012-b55f-bdc916deac2d-4
│   ├── [ 769] manifest-list-4ad2bd0d-79be-4012-b55f-bdc916deac2d-5
│   ├── [ 980] manifest-list-4ad2bd0d-79be-4012-b55f-bdc916deac2d-6
│   ├── [ 769] manifest-list-4ad2bd0d-79be-4012-b55f-bdc916deac2d-7
│   ├── [ 989] manifest-list-4ad2bd0d-79be-4012-b55f-bdc916deac2d-8
│   ├── [ 769] manifest-list-4ad2bd0d-79be-4012-b55f-bdc916deac2d-9
│   ├── [ 676] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-0
│   ├── [ 769] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-1
│   ├── [ 807] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-10
│   ├── [ 769] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-11
│   ├── [ 815] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-12
│   ├── [ 769] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-13
│   ├── [ 835] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-14
│   ├── [ 769] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-15
│   ├── [ 840] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-16
│   ├── [ 769] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-17
│   ├── [ 848] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-18
│   ├── [ 769] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-19
│   ├── [ 769] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-2
│   ├── [ 857] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-20
│   ├── [ 779] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-21
│   ├── [ 869] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-22
│   ├── [ 779] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-23
│   ├── [ 881] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-24
│   ├── [ 779] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-25
│   ├── [ 892] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-26
│   ├── [ 779] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-27
│   ├── [ 901] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-28
│   ├── [ 779] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-29
│   ├── [ 769] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-3
│   ├── [ 781] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-4
│   ├── [ 769] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-5
│   ├── [ 790] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-6
│   ├── [ 769] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-7
│   ├── [ 798] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-8
│   └── [ 769] manifest-list-a7d7473b-5da2-492f-b686-ffb248643e96-9
├── [ 96] schema
│   └── [ 603] schema-0
└── [ 864] snapshot
├── [ 2] EARLIEST
├── [ 2] LATEST
├── [ 529] snapshot-1
├── [ 533] snapshot-10
├── [ 534] snapshot-11
├── [ 534] snapshot-12
├── [ 534] snapshot-13
├── [ 534] snapshot-14
├── [ 534] snapshot-15
├── [ 532] snapshot-16
├── [ 533] snapshot-17
├── [ 533] snapshot-18
├── [ 533] snapshot-19
├── [ 529] snapshot-2
├── [ 533] snapshot-20
├── [ 536] snapshot-21
├── [ 535] snapshot-22
├── [ 535] snapshot-23
├── [ 530] snapshot-3
├── [ 530] snapshot-4
├── [ 530] snapshot-5
├── [ 532] snapshot-6
├── [ 533] snapshot-7
├── [ 532] snapshot-8
└── [ 532] snapshot-9

Flink SQL支持增量窗口TopN

假设定义数据源:

字段 类型 备注
category String 股票分类
field String 股票id
ftime String 字符串类型时间,格式为yyyyMMddHHmmss

样例数据:

category field ftime
k stock_id-1 20220930164620
k stock_id-0 20220930164625
k stock_id-1 20220930164625
k stock_id-0 20220930164625
k stock_id-1 20220930164625
k stock_id-0 20220930164625
k stock_id-0 20220930164625
k stock_id-0 20220930164625
k stock_id-0 20220930164625
k stock_id-0 20220930164625
k stock_id-1 20220930164630
k stock_id-1 20220930164630
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
select 
fixedTime(stat_time, 'HH:mm:ss') as time,
category,
field,
pv,
rk
from (
SELECT category,
field,
stat_time,
pv,
row_number() over (partition by category,INCREMENT(stat_time, INTERVAL '5' SECOND, INTERVAL '5' SECOND, Time '00:00:00', true) order by pv desc) as rk
FROM (SELECT category,
field,
INCREMENT_SEGMENT_ROWTIME(rtime, INTERVAL '1' DAY, INTERVAL '5' SECOND, Time '16:00:00', true) as stat_time,
count(1) as pv
FROM PageViews
GROUP BY category,
field,
INCREMENT(rtime, INTERVAL '1' DAY, INTERVAL '5' SECOND, Time '16:00:00', true)
)
)
where rk <= 10

计算结果:

row_kind time category field pv rk
true 16:49:54 k stock_id-1 2 1
true 16:49:54 k stock_id-0 1 2
true 16:49:59 k stock_id-1 8 1
true 16:49:59 k stock_id-0 5 2
true 16:50:04 k stock_id-1 12 1
true 16:50:04 k stock_id-0 11 2
true 16:50:09 k stock_id-0 19 1
true 16:50:09 k stock_id-1 14 2
true 16:50:14 k stock_id-0 23 1
true 16:50:14 k stock_id-1 20 2
true 16:50:19 k stock_id-1 27 1
true 16:50:19 k stock_id-0 26 2
1
2
3
4
5
6
7
8
9
Plan after converting SqlNode to RelNode:

LogicalProject(EXPR$0=[FixedTime($2, _UTF-16LE'HH:mm:ss')], category=[$0], field=[$1], pv=[$3], rk=[$4])
LogicalFilter(condition=[<=($4, 10)])
LogicalProject(category=[$0], field=[$1], stat_time=[$2], pv=[$3], rk=[ROW_NUMBER() OVER (PARTITION BY $0, INCREMENT($2, 5000:INTERVAL SECOND, 5000:INTERVAL SECOND, 00:00:00, true) ORDER BY $3 DESC)])
LogicalProject(category=[$0], field=[$1], stat_time=[INCREMENT_SEGMENT_ROWTIME($2)], pv=[$3])
LogicalAggregate(group=[{0, 1, 2}], pv=[COUNT()])
LogicalProject(category=[$0], field=[$1], rtime=[INCREMENT($2, 86400000:INTERVAL DAY, 5000:INTERVAL SECOND, 16:00:00, true)], $f3=[1])
FlinkLogicalDataStreamScan(id=[2], fields=[category, field, rtime])

Source 里面有三个字段: category, field, rtime
其中的rtime字段是内置字段rtime,来标示事件时间。通过tableEnv.registerDataStream("PageViews", stockClickStream, "category, field, rtime.rowtime"); 指定。

1
2
3
4
5
6
7
8
9
10
org.apache.calcite.plan.RelOptPlanner org.apache.calcite.plan.hep.HepPlanner.dumpGraph(HepPlanner.java:1027)- 
Breadth-first from root: {
rel#20:HepRelVertex#20 = rel#19:LogicalProject.NONE(input=HepRelVertex#18,exprs=[_UTF-16LE'Top10', FixedTime($2, _UTF-16LE'yyyyMMdd-HHmmss'), $0, $1, $3, $4]), rowcount=5.0, cumulative cost={241.25 rows, 630.0 cpu, 3600.0 io}
rel#18:HepRelVertex#18 = rel#17:LogicalFilter.NONE(input=HepRelVertex#16,condition=<=($4, 10)), rowcount=5.0, cumulative cost={236.25 rows, 600.0 cpu, 3600.0 io}
rel#16:HepRelVertex#16 = rel#15:LogicalProject.NONE(input=HepRelVertex#14,inputs=0..3,exprs=[ROW_NUMBER() OVER (PARTITION BY $0, INCREMENT($2, 5000:INTERVAL SECOND, 5000:INTERVAL SECOND, 00:00:00, true) ORDER BY $3 DESC)]), rowcount=10.0, cumulative cost={231.25 rows, 590.0 cpu, 3600.0 io}
rel#14:HepRelVertex#14 = rel#13:LogicalProject.NONE(input=HepRelVertex#12,inputs=0..1,exprs=[INCREMENT_SEGMENT_ROWTIME($2), $3]), rowcount=10.0, cumulative cost={221.25 rows, 540.0 cpu, 3600.0 io}
rel#12:HepRelVertex#12 = rel#11:LogicalAggregate.NONE(input=HepRelVertex#10,group={0, 1, 2},pv=COUNT()), rowcount=10.0, cumulative cost={211.25 rows, 500.0 cpu, 3600.0 io}
rel#10:HepRelVertex#10 = rel#9:LogicalProject.NONE(input=HepRelVertex#8,inputs=0..1,exprs=[INCREMENT($2, 86400000:INTERVAL DAY, 5000:INTERVAL SECOND, 16:00:00, true), 1]), rowcount=100.0, cumulative cost={200.0 rows, 500.0 cpu, 3600.0 io}
rel#8:HepRelVertex#8 = rel#1:FlinkLogicalDataStreamScan.LOGICAL(id=2,fields=page, event, rtime), rowcount=100.0, cumulative cost={100.0 rows, 100.0 cpu, 3600.0 io}
}

cumulative cost: 累计代价

参考 最强最全面的大数据 SQL 面试题和答案

每个学生最好成绩的科目

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
CREATE TABLE `course` (
`student` varchar(20) DEFAULT NULL,
`course` varchar(20) DEFAULT NULL,
`score` int DEFAULT NULL
) ENGINE = InnoDB CHARSET = utf8mb4 COLLATE utf8mb4_0900_ai_ci;
insert into `course` values
('王五','语文','89'),
('王五','数学','136'),
('王五','理综','130'),
('王五','英语','92'),
('张三','英语','110'),
('张三','数学','110'),
('张三','语文','110'),
('张三','理综','110'),
('李四','理综','109'),
('李四','语文','120'),
('李四','数学','100'),
('李四','英语','120');
1
select * from course;
student course score
王五 语文 89
王五 数学 136
王五 理综 130
王五 英语 92
张三 英语 110
张三 数学 110
张三 语文 110
张三 理综 110
李四 理综 109
李四 语文 120
李四 数学 100
李四 英语 120

求每个学生成绩最好的科目

可能成绩最好的科目,分数一样

1
2
3
4
5
6
7
8
SELECT m.student, m.max_score, c.course
FROM (
SELECT student, MAX(score) AS max_score
FROM course
GROUP BY student
) m, course c
WHERE m.student = c.student
AND m.max_score = c.score;
student max_score course
王五 136 数学
张三 110 英语
张三 110 数学
张三 110 语文
张三 110 理综
李四 120 语文
李四 120 英语

行列转换

样例数据准备

1
2
3
4
5
6
7
8
9
create database sql_case;
-- 年份-部门-绩效表%%
CREATE TABLE `t1` (
`fyear` int DEFAULT NULL COMMENT '年份',
`fdept` char(2) DEFAULT NULL COMMENT '部门',
`fscore` int DEFAULT NULL COMMENT '绩效评分'
) ENGINE = InnoDB CHARSET = utf8mb4 COLLATE utf8mb4_0900_ai_ci

insert into t1 values ('2014', 'B', '9'), ('2015', 'A', '8'), ('2014', 'A', '10'), ('2015', 'B', '7');

多行转多列

1
select * from t1
fyear fdept fscore
2014 B 9
2015 A 8
2014 A 10
2015 B 7

求每年各部门的分数

1
2
3
4
5
6
7
8
9
10
SELECT 
fyear,
MAX(CASE
WHEN fdept = 'a' THEN fscore
END) as fdept_a,
MAX(CASE
WHEN fdept = 'b' THEN fscore
END) as fdept_b
FROM t1
GROUP BY fyear;
fyear fdept_a fdept_b
2014 10 9
2015 8 7

one more thing

错误写法

1
2
3
4
5
6
7
8
9
10
SELECT 
fyear,
CASE
WHEN fdept = 'a' THEN fscore
END as fdept_a,
CASE
WHEN fdept = 'b' THEN fscore
END as fdept_b
FROM t1
GROUP BY fyear;

报错

1
(1055, "Expression #2 of SELECT list is not in GROUP BY clause and contains nonaggregated column 'sql_case.t1.fdept' which is not functionally dependent on columns in GROUP BY clause; this is incompatible with sql_mode=only_full_group_by")

参考 https://www.cnblogs.com/jim2016/p/6322703.html

对于GROUP BY聚合操作,如果在SELECT中的列,没有在GROUP BY中出现,那么这个SQL是不合法的,因为列不在GROUP BY从句中,也就是说查出来的列必须在group by后面出现否则就会报错,或者这个字段出现在聚合函数里面。

只选择出现在group by后面的列,或者给列增加聚合函数

如何关闭ONLY_FULL_GROUP_BY选项?

服务级别(重启Mysql服务):

1
2
3
set @@GLOBAL.sql_mode='';

set sql_mode ='STRICT_TRANS_TABLES,NO_ZERO_IN_DATE,NO_ZERO_DATE,ERROR_FOR_DIVISION_BY_ZERO,NO_AUTO_CREATE_USER,NO_ENGINE_SUBSTITUTION';

全局配置:

在 [mysqld]和[mysql]下添加

1
sql_mode ='STRICT_TRANS_TABLES,NO_ZERO_IN_DATE,NO_ZERO_DATE,ERROR_FOR_DIVISION_BY_ZERO,NO_AUTO_CREATE_USER,NO_ENGINE_SUBSTITUTION';

2. 多列转多行

问题描述:将问题一的结果转成源表,问题一结果表名为t1_2

1
2
3
4
5
6
7
8
9
10
11
create table t1_2
SELECT
fyear,
MAX(CASE
WHEN fdept = 'a' THEN fscore
END) as fdept_a,
MAX(CASE
WHEN fdept = 'b' THEN fscore
END) as fdept_b
FROM t1
GROUP BY fyear;
1
select * from t1_2;
fyear fdept_a fdept_b
2014 10 9
2015 8 7

答案

1
2
3
4
5
6
7
8
9
SELECT fyear, fdept, fscore
FROM (
SELECT fyear, 'A' AS fdept, fdept_a AS fscore
FROM t1_2
UNION ALL
SELECT fyear, 'B' AS fdept, fdept_b AS fscore
FROM t1_2
) t
-- 必须写别名,否则报错: (1248, 'Every derived table must have its own alias')
fyear fdept fscore
2014 A 10
2015 A 8
2014 B 9
2015 B 7

3. 同一部门有多个绩效情况下,多行转多列

2014年公司组织架构调整,导致部门出现多个绩效,业务及人员不同,无法合并算绩效

1
2
3
4
5
6
7
8
CREATE TABLE t1_3 LIKE t1;

INSERT INTO t1_3
VALUES ('2014', 'B', '9'),
('2015', 'A', '8'),
('2014', 'A', '10'),
('2015', 'B', '7'),
('2014', 'B', '6');
fyear fdept fscore
2014 B 9
2015 A 8
2014 A 10
2015 B 7
2014 B 6
1
2
3
4
5
6
SELECT 
fyear,
fdept,
group_concat(fscore order by fscore separator ',') as fscore
FROM t1_3
GROUP BY fyear,fdept;
fyear fdept fscore
2014 A 10
2014 B 9,6
2015 A 8
2015 B 7
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
SELECT 
fyear,
MAX(CASE
WHEN fdept = 'a' THEN fscore
END) as fdept_a,
MAX(CASE
WHEN fdept = 'b' THEN fscore
END) as fdept_b
FROM (
SELECT
fyear,
fdept,
group_concat(fscore order by fscore separator ',') as fscore
FROM t1_3
GROUP BY fyear,fdept
) t
GROUP BY fyear;
fyear fdept_a fdept_b
2014 10 6,9
2015 8 7

排名中取他值

1
2
3
4
5
6
7
8
create table t2 like t1;
-- init values
insert into t2 values
('2014', 'A', '3'),
('2014', 'B', '1'),
('2014', 'C', '2'),
('2015', 'A', '4'),
('2015', 'D', '3');
1
select * from t2;
fyear fdept fscore
2014 A 3
2014 B 1
2014 C 2
2015 A 4
2015 D 3

按a分组取b字段最小时对应的c字段

1
2
3
4
5
6
7
8
9
10
select 
fyear,
fdept,
fscore,
row_number() over
(
partition by fyear
order by fdept
) as rn
from t2;
fyear fdept fscore rn
2014 A 3 1
2014 B 1 2
2014 C 2 3
2015 A 4 1
2015 D 3 2
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
select 
fyear,
fscore as fscore_of_min_fdept
from
(
select
fyear,
fdept,
fscore,
row_number() over
(
partition by fyear
order by fdept
) as rn
from t2
) tt where tt.rn=1;
fyear fscore_of_min_fdept
2014 3
2015 4

按a分组取b字段排第二时对应的c字段

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
select 
fyear,
fscore as fscore_of_sec_min_fdept
from
(
select
fyear,
fdept,
fscore,
row_number() over
(
partition by fyear
order by fdept
) as rn
from t2
) tt where tt.rn=2;
fyear fscore_of_sec_min_fdept
2014 1
2015 3

按a分组取b字段最小和最大时对应的c字段

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
select 
fyear,
fdept,
fscore,
row_number() over
(
partition by fyear
order by fdept
) as rn_asc,
row_number() over
(
partition by fyear
order by fdept desc
) as rn_desc
from t2;
fyear fdept fscore rn_asc rn_desc
2014 C 2 3 1
2014 B 1 2 2
2014 A 3 1 3
2015 D 3 2 1
2015 A 4 1 2
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
select
fyear,
min(if(rn_asc = 1, fscore, null)) as fscore_of_min_dept,
max(if(rn_desc = 1, fscore, null)) as fscore_of_max_dept
from
(
select
fyear,
fdept,
fscore,
row_number() over
(
partition by fyear
order by fdept
) as rn_asc,
row_number() over
(
partition by fyear
order by fdept desc
) as rn_desc
from t2
) tt
where tt.rn_asc = 1 or tt.rn_desc =1
group by fyear
fyear fscore_of_min_dept fscore_of_max_dept
2014 3 2
2015 4 3

按a分组取b字段第二小和第二大时对应的c字段

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
select
fyear,
min(if(rn_asc = 2, fscore, null)) as score_of_second_min_dept,
max(if(rn_desc = 2, fscore, null)) as score_of_second_max_dept
from
(
select
fyear,
fdept,
fscore,
row_number() over
(
partition by fyear
order by fdept
) as rn_asc,
row_number() over
(
partition by fyear
order by fdept desc
) as rn_desc
from t2
) tt
where tt.rn_asc = 2 or tt.rn_desc =2
group by fyear
fyear score_of_second_min_dept score_of_second_max_dept
2014 1 1
2015 3 4

按a分组取b字段前两小和前两大时对应的c字段

需保持fdept字段最小、最大排首位

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
select
fyear as `year`,
group_concat(if(rn_asc <= 2, fscore, null) order by rn_asc) as score_of_less_2_min_dept,
group_concat(if(rn_desc <= 2, fscore, null) order by rn_desc) as score_of_less_2_max_dept
from
(
select
fyear,
fdept,
fscore,
row_number() over
(
partition by fyear
order by fdept
) as rn_asc,
row_number() over
(
partition by fyear
order by fdept desc
) as rn_desc
from t2
) tt
where tt.rn_asc <= 2 or tt.rn_desc <=2
group by fyear
year score_of_less_2_min_dept score_of_less_2_max_dept
2014 3,1 2,1
2015 4,3 3,4

累计求值

1
2
3
4
5
6
7
8
create table t3 like t1;

insert into t3 values
('2015', 'A', '4'),
('2014', 'A', '3'),
('2014', 'C', '2'),
('2014', 'B', '1'),
('2015', 'D', '3');
fyear fdept fscore
2015 A 4
2014 A 3
2014 C 2
2014 B 1
2015 D 3

按a分组按b字段排序,对c累计求和

1
2
3
4
5
6
7
8
9
10
select 
fyear,
fdept,
fscore,
sum(fscore) over
(
partition by fyear order by fdept
) as accu_sum_score
from
t3;
fyear fdept fscore accu_sum_score
2014 A 3 3
2014 B 1 4
2014 C 2 6
2015 A 4 4
2015 D 3 7

按a分组按b字段排序,对c取累计平均值

1
2
3
4
5
6
7
8
9
10
select 
fyear,
fdept,
fscore,
avg(fscore) over
(
partition by fyear order by fdept
) as accu_avg_score
from
t3;
fyear fdept fscore accu_avg_score
2014 A 3 3.0000
2014 B 1 2.0000
2014 C 2 2.0000
2015 A 4 4.0000
2015 D 3 3.5000

按a分组按b字段排序,对b取累计排名比例

1
2
3
4
5
6
7
select
fyear as `year`,
fdept as `dept`,
fscore as `score`,
round(row_number() over (partition by fyear order by fdept)/(count(fscore) over(partition by fyear)),2) as accu_avg_score
from t3
order by fyear,fdept
year dept score accu_avg_score
2014 A 3 0.33
2014 B 1 0.67
2014 C 2 1.00
2015 A 4 0.50
2015 D 3 1.00

按a分组按b字段排序,对b取累计求和比例

1
2
3
4
5
6
7
select
fyear as `year`,
fdept as `dept`,
fscore as `score`,
round(row_number() over (partition by fyear order by fdept)/(sum(fscore) over(partition by fyear)),2) as accu_ratio_score
from t3
order by fyear,fdept
year dept score accu_ratio_score
2014 A 3 0.17
2014 B 1 0.33
2014 C 2 0.50
2015 A 4 0.14
2015 D 3 0.29

窗口大小控制

1
2
3
4
5
6
7
8
create table t4 like t1;

insert into t4 values
('2014', 'A', '3'),
('2014', 'B', '1'),
('2014', 'C', '2'),
('2015', 'A', '4'),
('2015', 'D', '3');
fyear fdept fscore
2014 A 3
2014 B 1
2014 C 2
2015 A 4
2015 D 3

按a分组按b字段排序,对c取前后各一行的和

1
2
3
4
5
select 
fyear as `year`,
fdept as dept,
lag(fscore,1,0) over(partition by fyear order by fdept) | lead(fscore,1,0) over(partition by fyear order by fdept) as sum_range_score
from t4;
year dept sum_range_score
2014 A 1
2014 B 5
2014 C 1
2015 A 3
2015 D 4

按a分组按b字段排序,对c取平均值

前一行与当前行的均值

1
2
3
4
5
6
select
fyear as `year`,
fdept as `dept`,
fscore as `score`,
lag(fscore,1) over(partition by fyear order by fdept) as lag_score
from t4;
year dept score lag_score
2014 A 3
2014 B 1 3
2014 C 2 1
2015 A 4
2015 D 3 4
1
2
3
4
5
6
7
8
9
10
11
12
13
14
select
`year`,
`dept`,
`score`,
case when lag_score is null then score else (score | lag_score)/2 end as avg_lag2
from
(
select
fyear as `year`,
fdept as `dept`,
fscore as `score`,
lag(fscore,1) over(partition by fyear order by fdept) as lag_score
from t4
) tt
year dept score avg_lag2
2014 A 3 3
2014 B 1 2.0000
2014 C 2 1.5000
2015 A 4 4
2015 D 3 3.5000

产生连续数值

数据扩充与收缩

合并与拆分

模拟循环操作

不使用distinct或group by去重

容器–反转内容

多容器–成对提取数据

多容器–转多行

抽象分组–断点排序

业务逻辑的分类与抽象–时效

时间序列–进度及剩余

时间序列–构造日期

时间序列–构造累计日期

时间序列–构造连续日期

时间序列–取多个字段最新的值

时间序列–补全数据

时间序列–取最新完成状态的前一个状态

非等值连接–范围匹配

非等值连接–最新匹配

N指标–累计去重

concept 中文 object
Relational algebra 关系代数 RelNode
Row expressions 行表达式 RexNode
Traits RelTrait
Conventions Convention
Rules RelOptRule
Planners RelOptPlanner
Programs Program

关系代数

关系代数 对应的Spark操作
TableScan SparkTableScan
Project SparkProject
Filter SparkFilter
Aggregate SparkAggregate
Join SparkJoin
Union SparkUnion
Intersect SparkIntersect
Sort SparkSort