其他项目相关¶
写一个prom sql语句查询CPU和内存信息¶
查询CPU使用率:
# 查询近5分钟的CPU使用率
100 - (avg by(instance) (irate(node_cpu_seconds_total{mode="idle"}[5m])) * 100
查询内存使用率:
# 查询内存使用率
100 - (node_memory_MemAvailable_bytes / node_memory_MemTotal_bytes) * 100
简单说一下,第一个语句的意思:nodecpuseconds_total是CPU各个模式的累计时间,irate取最近5分钟的瞬时增长率,mode="idle"是空闲时间,用100减去空闲就是使用率,再按instance分组。第二个是可用内存除以总内存,用100减就是内存使用率。
问题 1:DolphinScheduler K8s 连接池优化项目 你当时做这个连接池优化,先讲下原始架构存在的核心痛点是什么?为什么原生逻辑没有复用客户端,而是每次新建? 另外你用 SHA-256 对 kubeconfig 生成集群唯一标识,有没有考虑过 kubeconfig 内容相同但集群实际不同的边界场景?是怎么规避的? 问题 2:大数据数仓项目 你搭建了实时 + 离线一体化电商数仓,Flink 消费 Kafka 做实时指标延迟控制在 5s 内,说下你在 Flink 侧做了哪些调优手段降低延迟? 离线层你用 Hive 分层建模,DWD/DWS/ADS 三层分别承载什么职责,举一个你落地的业务指标计算案例。 问题 3:基础 Java 底层问题 简历写你精通 Java8+、多线程、JVM,连接池本身就是典型多线程场景: 你写 K8s 连接池时,池内空闲客户端采用什么并发容器存储?为什么选择它而非其他容器? 连接池的空闲清理定时任务,如果清理线程和业务获取客户端线程并发冲突,你是怎么保证线程安全的? 问题 4:开源项目 Bigtop Manager 你在 openEuler 开源实习做指标采集,采集 CPU、磁盘等 15 项指标,当时怎么平衡采集精度和服务器性能开销?采集线程调度策略具体做了哪些优化,最终资源消耗下降 20% 是怎么测算验证的?
第一原始架构的核心问题是不断创建客户端,就会造成创建和销毁的资源开销大,然后在多任务调度的时候会出现短时间大量创建客户端的客户端爆炸压力,为什么选择新建就是为了解决Watch机制对客户端的独占性需求,也是为了在多集群的任务背景下解决问题,于是选择了按任务缓存客户端。没有,默认认知是config相关必然指向了实际不同的边界场景,如果确实需要考虑这个问题的话,可以在config里面添加更详细的内容进行区分
第二,# Flink消费Kafka 5s低延迟核心调优(面试精简版)
为保障Flink消费Kafka、实时计算写入Doris全链路延迟稳定控制在5s以内,我主要从源头消费、Flink内核、业务算子、状态检查点、Sink下游、底层优化六个核心维度做调优,都是落地有效的关键手段,具体核心要点如下:
一、Kafka Source 源头低延迟调优(核心)¶
核心思路:放弃高吞吐批量消费,优先保障低延迟,消除数据拉取等待。 - 并行度对齐:严格保证 Flink Source 并行度 = Kafka Topic 分区数,一对一独占消费,杜绝单线程多分区阻塞、空闲资源浪费。 - 消费者参数调优:核心修改两个参数,fetch.min.bytes=1(有一条数据就拉取,不攒批)、fetch.max.wait.ms=100(最大等待100ms,无数据立即返回),彻底消除拉取阻塞。 - Offset提交优化:关闭Kafka定时自动提交,改为仅Checkpoint提交Offset,减少网络IO开销,同时依托Flink状态保证精准消费。 - 快速重平衡:调低消费者会话超时时间,避免重平衡耗时过长导致数据堆积。
二、Flink 运行时核心调优¶
- 全链路并行度一致:Source、中间算子、Doris Sink 并行度统一,彻底消除上下游数据传输背压瓶颈。
- 开启流式调度:禁用批式调度,开启 Pipelined 流式调度,数据逐条实时下发,无需等待上游批次完成。
- 最大化算子链:无shuffle的窄依赖算子(map、filter)全部合并链式执行,避免线程切换、序列化反序列化开销。
- 资源隔离:实时任务单独配置slot分组,不与离线任务共享资源,防止资源抢占导致延迟飙升。
三、业务算子与窗口核心优化(关键)¶
针对电商实时数仓聚合场景,重点优化窗口和乱序数据处理:
- 轻量化秒级窗口:摒弃大窗口,使用1s滚动窗口+增量聚合ReduceFunction,每条数据实时更新聚合结果,每秒触发计算输出,不缓存全量窗口数据,拒绝延迟堆积。
- 极致水位线优化:设置乱序容忍500ms,适配业务轻微乱序场景;开启空闲分区检测,避免分区无数据导致全局水位停滞、窗口无法触发。
- 热点Key优化:对品牌、渠道等电商热点Key做加盐打散+本地预聚合,解决单线程计算瓶颈。
- 前置数据过滤:提前过滤空值、无效测试数据,减少下游无效计算压力。
四、Checkpoint \& 状态后端调优(解决延迟抖动)¶
Checkpoint卡顿是延迟抖动的核心原因,核心优化为轻量、异步快照: - 高频轻量快照:设置1s一次Checkpoint,单次快照数据量小、IO耗时短,避免长间隔快照故障重放延迟。 - RocksDB增量异步快照:开启增量快照仅持久化变化状态,搭配异步写入,快照操作不阻塞业务数据计算。 - 状态自动清理:配置24h状态TTL过期清理,避免状态无限膨胀导致GC、快照变慢。
五、Doris Sink 下游调优(杜绝下游背压)¶
Flink写入Doris是最常见延迟瓶颈,核心优化小批量高频写入: - 调低刷写阈值:设置batchSize=100、batchIntervalMs=500,满足任意条件立即刷写,不缓存大量数据。 - 异步并行写入:开启Sink多线程异步提交,提升写入吞吐。 - Doris侧适配:合理分片分区、扩容BE节点,保证数据库写入低响应延迟,不反压上游Flink。
六、底层兜底优化¶
- GC调优:使用G1垃圾收集器,限制单次GC停顿10ms内,避免长时间STW卡顿。
- 网络优化:关闭TCP Nagle算法,小包即时发送,消除网络传输延迟。
- 实时监控兜底:监控上下游读写数据量,实时感知背压,异常及时告警。
一、数仓三层核心职责(面试精简版)¶
1. DWD 明细层¶
核心定位:清洗标准化原始数据,存储最细粒度业务明细,是数仓底层基础。
核心职责:清洗脏数据、拆解嵌套字段、统一数据粒度、维度退化,区分增量/全量快照表。
存储规范:ORC+Snappy压缩,按日期分区。
示例:dwd_trade_order_detail_inc,一行对应订单单个商品明细。
2. DWS 汇总宽表层¶
核心定位:基于DWD明细,做轻度预聚合、宽表化整合,实现指标复用,提升查询效率。
核心职责:多维度预聚合、合并常用维度为宽表,规避上层重复计算,承接通用派生指标。
存储规范:ORC+Snappy,按日期分区。
示例:dws_trade_province_category_day_agg,日-省份-品类维度交易汇总。
3. ADS 应用报表层¶
核心定位:面向业务和BI可视化,产出最终业务指标,直接落地报表展示。
核心职责:统一业务口径、计算环比/占比/转化率等衍生指标,适配可视化数据源。
存储规范:Hive加工后同步MySQL,供Superset使用。
示例:ads_order_by_province,各省日订单统计报表。
二、落地业务指标案例(精简面试版)¶
1. 业务需求¶
统计每日各类型优惠券领取量、使用量、抵扣金额、使用率及占比,用于Superset图表展示。
2. 分层落地流程¶
① DWD层(明细清洗):读取ODS优惠券日志,清洗脏数据、拆解JSON、关联维度表退化券类型名称,生成dwd_tool_coupon_use_inc明细事实表。
② DWS层(预聚合):按日期、优惠券类型分组,预聚合领取数、使用数、抵扣金额,生成日维度汇总宽表,复用基础指标。
③ ADS层(指标计算):基于DWS数据,计算优惠券使用率、金额占比等衍生指标,产出最终报表数据。
④ 可视化落地:将ADS表同步至MySQL,对接Superset制作趋势卡片、占比饼图,配置自动刷新仪表盘。
三、三层核心总结(面试速记)¶
| 分层 | 粒度 | 核心作用 | 用途 |
|---|---|---|---|
| DWD | 最细明细 | 数据清洗、标准化 | 底层明细支撑 |
| DWS | 维度聚合 | 预聚合、指标复用 | 减少重复计算 |
| ADS | 业务报表 | 衍生指标、口径对齐 | 业务报表、BI展示 |
四、Hive存储格式选择规则(面试精简版)¶
核心原则:数仓分层不同,读写场景不同,对应选择不同存储格式,离线数仓主流优先ORC,日志原始数据可选Parquet/Text。
1. ODS层(原始数据层)¶
场景:存储原始日志、业务binlog数据,多为增量写入、偶尔全量回溯,极少聚合查询。
选择:Text文本格式 + Gzip压缩
理由:兼容性最强,采集接入无需序列化,写入速度快;缺点是不支持索引、查询慢,符合ODS只存不查的定位。
2. DWD/DWS层(明细/汇总层)¶
场景:频繁过滤、聚合、关联查询,数据读写量大、计算频次高,对查询性能和压缩比要求高。
选择:ORC列式存储 + Snappy压缩
理由:自带索引、谓词下推,过滤查询效率极高;列式存储适配聚合计算,压缩比高、节省存储,完美适配明细和汇总层高频计算场景。
3. ADS层(报表应用层)¶
场景:数据量小、查询频次高,多为整表读取用于报表展示、数据同步。
选择:ORC + Snappy(沿用上层),无需特殊格式
理由:统一数仓格式规范,读写稳定,方便直接同步至MySQL供Superset可视化使用。
第三,从 PR 的 KubernetesClientPool.java 实现看,空闲客户端存储采用的是 LinkedBlockingQueue<PooledClient>。具体理由与选型分析如下:
一、选型依据:为什么是 LinkedBlockingQueue?¶
| 特性需求 | LinkedBlockingQueue 的匹配点 |
|---|---|
| 支持超时等待 | 它实现了 BlockingQueue 接口,原生支持 poll(long timeout, TimeUnit unit),这正好对应 borrowObject() 里“池满时等待 maxWaitMs”的需求 |
| 适合“先入先出”的借还模型 | 队列的 FIFO 特性与连接池“先还先借”的朴素公平策略匹配,避免个别连接过度使用/闲置不均 |
| 高并发入队出队友好 | LinkedBlockingQueue 内部用了两把锁分离(takeLock + putLock):入队和出队操作分别加锁,读少写多时性能比 ArrayBlockingQueue 的单锁更好 |
| 动态容量适配 | 默认是无界(或可设 capacity)的,能自然适配 minIdle 到 maxSize 的动态变化,不用预先分配固定大数组 |
二、为什么没选其他常见容器?¶
1. 没选 ConcurrentLinkedQueue¶
- 它是非阻塞的,没有
poll(timeout),只能自旋或自己写wait/notify - 连接池在满负载下需要“等待释放”的语义,
BlockingQueue更直接
2. 没选 ArrayBlockingQueue¶
- 它是单锁模型,入队出队共用一把锁,在并发较高的“归还-借用”场景下,锁竞争比
LinkedBlockingQueue更大 - 容量固定,无法灵活应对
minIdle/maxSize的动态变更(虽然本 PR 没做运行时配置热更,但设计上预留了弹性)
3. 没选 ConcurrentLinkedDeque / LinkedBlockingDeque¶
- 连接池不需要“双端”操作(不需要从头部/尾部同时取/放),用
Deque属于过度设计 - 没有比
LinkedBlockingQueue带来额外收益
三、一个值得注意的“双重同步”细节¶
虽然 idleClients 用了 LinkedBlockingQueue(本身线程安全),但 borrowObject、returnObject、cleanupIdle 三个方法都额外加了 synchronized:
public synchronized KubernetesClient borrowObject() throws Exception {
PooledClient client = idleClients.poll();
// ...
}
public synchronized void returnObject(KubernetesClient client) {
// ...
}
原因:这三个方法除了操作 idleClients,还同时操作 activeClients(一个普通 HashSet)和 createdCount(虽然是 AtomicInteger,但与其他操作需要一起原子化),所以需要用 synchronized 保护整体不变性。
这种“并发容器 + 外层同步”的做法是合理的,因为我们要保证的是「idleClients、activeClients、createdCount 三者状态一致性」,而不仅仅是 idleClients 本身的线程安全。
如果要进一步优化,可考虑:
- 把 activeClients 改成 ConcurrentHashMap<PooledClient, Boolean>(或用 ConcurrentHashMap.newKeySet())
- 减少 synchronized 的粒度,只在必要的复合操作上加锁
从该 PR 的实现看,清理线程与业务线程的并发安全是通过「对集群池实例加同一把 synchronized 锁」来保证的。具体分析如下:
一、核心安全机制:所有修改池状态的方法共用 ClusterClientPool.this 锁¶
ClusterClientPool 中三个会操作 idleClients、activeClients、createdCount 的方法都加了 synchronized:
public synchronized KubernetesClient borrowObject() throws Exception { ... }
public synchronized void returnObject(KubernetesClient client) { ... }
public synchronized void cleanupIdle() { ... }
效果:
- 同一时刻,只能有一个线程(要么是业务借/还线程,要么是清理线程)进入这三个方法之一
- 保证了 idleClients、activeClients、createdCount 三者的状态一致性,不会出现“业务刚把连接放回 idle,清理线程就把它拿走并关闭”的竞态条件
二、cleanupIdle 内部的额外安全处理¶
即使加了 synchronized,cleanupIdle 内部实现仍做了一个安全细节:
public synchronized void cleanupIdle() {
long now = System.currentTimeMillis();
// 先转成数组,再遍历
PooledClient[] clients = idleClients.toArray(new PooledClient[0]);
int keepIdle = Math.max(config.getMinIdle(), 0);
int removeCount = 0;
for (PooledClient client : clients) {
if (idleClients.size() - removeCount > keepIdle
&& now - client.lastUsedTime > config.getIdleTimeoutMs()) {
if (idleClients.remove(client)) {
closeClient(client);
removeCount++;
}
}
}
}
为什么先 toArray 再遍历?
- 虽然方法整体加了锁,但 idleClients.remove(client) 会修改队列结构
- 如果直接在 BlockingQueue 上做 for-each 遍历,遍历时删除元素可能导致 ConcurrentModificationException 或漏掉元素
- 先转成数组相当于“拍了快照”,然后基于快照判断超时,再用 idleClients.remove(client) 从原队列中精确移除(因为此时仍在锁内,所以 remove 是安全的)
三、清理线程的启动与守护线程设置¶
private void startCleanupThread() {
Thread cleanupThread = new Thread(() -> {
while (true) {
try {
Thread.sleep(30000);
cleanupIdleClients();
} catch (InterruptedException e) {
log.warn("Cleanup thread interrupted", e);
Thread.currentThread().interrupt();
break;
}
}
}, "k8s-client-cleanup-thread");
cleanupThread.setDaemon(true);
cleanupThread.start();
}
- 设为
daemon线程,保证 JVM 退出时不会被该线程阻塞 - 正确处理了
InterruptedException:恢复中断位并退出,避免吞掉中断导致线程无法正常结束
四、有没有锁粒度优化空间?¶
当前做法是对整个 ClusterClientPool 实例加锁,优点是实现简单、状态一致性一目了然;缺点是在清理时,业务线程借还客户端会被短暂阻塞。
如果后续要优化锁粒度,可以考虑:
1. 用 ReadWriteLock:
- 借还操作用 writeLock(因为会修改 idle/active)
- 清理操作也用 writeLock(因为会清理 idle)
- 本质上和 synchronized 区别不大,因为清理也是写操作
2. 用 ConcurrentHashMap.newKeySet() 替代 HashSet<PooledClient> 做 activeClients,并把 idleClients 的 LinkedBlockingQueue 并发能力用起来,减少 synchronized 的包裹范围
3. 把“判断超时与关闭”移到锁外,只在“从 idle 队列移除”那一步加锁
不过对于 DolphinScheduler 的场景(K8s 任务量级通常不会特别极端,清理间隔 30 秒),当前的粗粒度 synchronized 是“简单且足够安全”的选择,没有过度设计的必要。