GPT5.2回答大数据开发工程师应该具备的能力
下面是一份“面向大数据开发工程师”的 Java 学习笔记(更偏工程落地 + 性能 + JVM + 并发 + 生态集成),学完后基本能覆盖你在 Hadoop/Spark/Flink/Kafka/ES/OLAP 等场景里写生产级 Java 的核心需求。你可以把它当作:路线图 + 速查笔记 + 易错点清单 + 练习清单。
0. 学习目标(大数据工程视角)
你需要的 Java 能力通常不是“会语法就行”,而是:
- 写得对:类型系统、集合、序列化、并发安全、资源释放、异常边界清晰
- 跑得稳:JVM 内存/GC/类加载、线程模型、反压与超时、日志与指标
- 跑得快:对象分配、装箱拆箱、集合选择、I/O 模型、锁竞争、热点优化
- 好交付:Maven/Gradle、依赖冲突、shade、可测试性、可观测性
- 能对接生态:Kafka 客户端、HDFS、YARN、Spark/Flink connector、序列化格式(Avro/Protobuf/JSON/Parquet/ORC)
1. 基础语法(只记工程里真的会踩的点)
1.1 基本类型与装箱
- 基本类型:
int/long/double/boolean等;引用类型:Integer/Long/... - 大数据常见坑:频繁装箱导致 GC 压力
- 例:
List<Long>海量数据比LongStream/long[]更容易炸内存
- 例:
Integer缓存:-128~127(实现细节),别用==比较包装类型,统一用equals
1.2 字符串与编码
String不可变;拼接循环用StringBuilder- 编码:所有 I/O 明确 charset(UTF-8)
String.intern()谨慎,容易把压力转到常量池/元空间
1.3 控制流与可读性
- 条件/循环不难,关键是:避免过深嵌套、提早返回、明确边界条件
- 生产里更推荐:小函数 + 清晰命名 > “聪明代码”
2. 面向对象与接口设计(写可维护代码的关键)
2.1 设计原则(够用版)
- 单一职责、依赖倒置(面向接口)
- 组合优于继承(大数据项目依赖多、类层次容易失控)
- 不要滥用
static全局状态(尤其在分布式任务/多线程里)
2.2 equals/hashCode/toString(集合行为核心)
- 放进
HashMap/HashSet的对象必须保证:equals一致性hashCode与equals兼容
- 常见坑:对象字段可变导致 hash 变化 → Map 取不出来
3. 泛型(你写框架/工具库时会救命)
3.1 核心概念
- 类型擦除:运行时没有泛型实参
- 通配符:
? extends T:只读(生产者)? super T:只写(消费者)
- Big Data 里常用于:通用转换器、反序列化器、通用 DAO/Client 封装
3.2 常见坑
List<?>基本不能add(除了null)- 泛型数组不安全:
new T[]不行,通常用List<T>或反射规避
4. 集合框架(大数据 Java 的基本功)
4.1 选型速查(高频)
ArrayList:随机读快,尾部追加快;扩容会拷贝LinkedList:多用于队列/双端队列(但一般ArrayDeque更好)HashMap:最常用;注意初始容量、负载因子;避免 key 可变TreeMap:有序,O(logn)ConcurrentHashMap:并发 Map 首选;别用Collections.synchronizedMap兜底ArrayDeque:高性能队列/栈,优于Stack/LinkedList
4.2 性能与内存
- 预估容量:大量 put 时设置初始容量避免 rehash
- 遍历删除:用迭代器
iterator.remove(),不要 for-each 里直接remove
5. 异常与资源管理(稳定性直接相关)
5.1 异常策略(建议你形成团队规范)
- 参数错误:
IllegalArgumentException - 状态错误:
IllegalStateException - 外部系统失败:自定义
XxxClientException包装原始异常(保留 cause) - 不要吞异常;日志要带上下文(topic、partition、offset、path、jobId 等)
5.2 try-with-resources(强烈推荐)
- 对所有
Closeable/AutoCloseable使用它(流、channel、client 等) - 分布式任务里资源泄漏会放大成“节点不稳定/FD 耗尽”
6. I/O 与 NIO(读写 HDFS/本地/网络都离不开)
6.1 I/O 模型
- BIO:阻塞流(简单但吞吐受限)
- NIO:Channel/Buffer(更接近高性能网络/文件 I/O)
- 你需要至少理解:
ByteBuffer的position/limit/flip/clear- 零拷贝相关概念(sendfile、mmap 是加分项)
6.2 实战建议
- 大文件:优先流式处理,避免一次性读入内存
- 压缩:gzip/snappy/lz4 不同 CPU/压缩率权衡(大数据里非常现实)
7. 并发与多线程(大数据工程的分水岭)
7.1 线程安全三件套
- 可见性:
volatile、锁、原子类 - 原子性:锁、CAS(Atomic)
- 有序性:happens-before(别硬背,理解 volatile/锁的语义)
7.2 线程池(必须掌握)
- 用
ThreadPoolExecutor明确参数:- core/max、队列、拒绝策略、线程命名
- 大数据常见坑:
- 不设队列边界 → 内存被任务堆爆
Executors.newFixedThreadPool默认无界队列(风险)
- 拒绝策略要与业务匹配:丢弃/CallerRuns/抛异常/记录报警
7.3 并发工具
ConcurrentHashMap:并发缓存/去重CountDownLatch:等待多任务完成Semaphore:限流CompletableFuture:异步编排(超时、异常合并很有用)ThreadLocal:谨慎(线程池复用会造成脏数据/内存泄漏;用完 remove)
8. JVM(大数据现场排障必备)
8.1 内存结构(够用版)
- 堆:新生代/老年代(不同 GC 略有差异)
- 元空间:类元数据(依赖冲突/动态代理多时关注)
- 直接内存:NIO/Netty 常用(堆外 OOM 常见)
8.2 GC 你要会看什么
你至少要能回答:
- 为什么 Full GC?(老年代、元空间、晋升失败、分配失败)
- GC 暂停是否影响 SLA?
- 对象为何大量分配?(装箱、字符串、集合扩容、日志拼接)
8.3 排障工具清单(常用)
jps/jstack/jmap/jcmd/jstat- heap dump + MAT/VisualVM/YourKit(定位内存泄漏)
- GC log(必开,生产排障靠它)
9. 序列化与数据格式(大数据“吞吐与成本”的核心)
9.1 Java 原生序列化(不推荐)
Serializable通常不建议用于跨版本、跨语言、高性能场景- 风险:性能差、易出安全问题、兼容性难控
9.2 推荐路线
- JSON:调试友好但体积大(注意字段命名、null、时间格式)
- Avro/Protobuf:强 schema、更高效(Kafka/Spark/Flink 常见)
- Parquet/ORC:列存(离线数仓主力),理解 schema 演进/压缩/字典编码更加分
10. 构建与依赖(Maven/Gradle + 依赖地狱处理)
10.1 你必须会的 Maven 点
dependency:tree看冲突scope:compile/provided/test/runtimeshade/relocation:打 fat-jar,避免与集群自带依赖冲突(Spark/Flink 场景尤其常见)- 统一版本管理:
dependencyManagement
10.2 典型冲突场景
- Hadoop/Spark 自带 guava、jackson、netty 版本与你项目不一致
→ 轻则告警,重则NoSuchMethodError/ClassNotFoundException
11. 日志、指标、可观测性(生产必备)
11.1 日志
- slf4j 门面 + logback/log4j2 实现
- 日志字段化:jobId、traceId、partition、offset、latency、retryCount
- 避免热路径字符串拼接:用参数化日志
11.2 Metrics
- Micrometer/Dropwizard Metrics(计数、耗时、直方图)
- 大数据任务最该打的指标:
- 吞吐(records/s、bytes/s)
- 延迟(p95/p99)
- 失败率、重试次数
- backlog/队列长度
- GC/堆使用、线程池队列
12. 测试(让你改代码更有底气)
- 单测:JUnit5 + AssertJ/Hamcrest
- Mock:Mockito(外部系统用 mock,别真连 Kafka/HDFS)
- 集成测试:Testcontainers(Kafka、Postgres 等)
- 属性测试/随机测试:对序列化、解析、边界条件很有帮助(加分项)
13. 网络编程与客户端(Kafka/HTTP/HBase 等常见)
- HTTP Client:连接池、超时(连接/读/写)、重试策略(幂等性)
- Kafka client:理解
- consumer group、rebalance
- offset 提交语义(at-least-once/exactly-once 相关)
- 背压与拉取批次(max.poll.records 等)
14. 安全与稳定性(别等事故教你)
- 反序列化安全:不要反序列化不可信数据(尤其 Java 原生序列化)
- 输入校验:路径、SQL、JSON 字段
- 资源上限:线程池队列、bulk 请求大小、buffer 上限、超时必设
- 降级与熔断:重试要有上限 + 指数退避 + 抖动
15. 大数据框架中的 Java“必懂交互点”
15.1 Spark(Java 视角)
- 闭包序列化:Java lambda 捕获外部变量 → 容易序列化失败/带入大对象
- 广播变量与累加器:避免在 executor 端频繁拉取大对象
- UDF:注意空值、类型、性能(避免频繁创建对象)
15.2 Flink(Java 视角)
- Operator 生命周期:open/close、状态后端、checkpoint 语义
- Watermark/时间语义:事件时间处理的坑比语法多
- 序列化:POJO 规则、Kryo、TypeInformation(性能差异很真实)
15.3 Hadoop/HDFS
- InputFormat/OutputFormat、Writable(历史遗留但仍会遇到)
- 小文件问题:合并、压缩、合理分区
16. 一份“学完就能上手”的练习清单(建议按顺序做)
- 实现一个高性能日志解析器:支持 gzip 输入、逐行处理、输出统计指标
- 写一个 Kafka Consumer 服务:手动提交 offset + 重试 + 死信队列(DLQ)
- 实现一个并发批处理器:线程池 + 有界队列 + backpressure + 指标
- 做一次 JVM 排障演练:故意制造内存泄漏/线程死锁,然后用工具定位
- 打包一个可在 Spark/Flink 集群跑的 fat-jar:解决一次依赖冲突(shade relocation)
- 实现 Avro/Protobuf schema 演进 demo:验证向后/向前兼容
17. 面试/自检速查(你可以当 checklist)
HashMap原理、扩容、为什么 key 不能乱变- 线程池 7 参数、为什么无界队列危险
volatile能解决什么不能解决什么- JVM:堆/元空间/直接内存,OOM 类型与排查思路
- GC 日志能看懂哪些关键信号
- 序列化方案对比:JSON vs Avro/Protobuf vs Parquet
- Maven 依赖冲突定位与解决
- 日志与 metrics 如何设计(你会比只会写业务的人强很多)
推荐学习顺序(最省时间的路径)
- 集合/泛型/异常/资源管理(打底)
- 并发 + 线程池(决定上限)
- JVM + GC + 工具(决定你能不能救火)
- I/O + 序列化 + 构建依赖(决定吞吐与交付)
- 框架对接与工程化(Spark/Flink/Kafka 实战)
下面继续把这份笔记往“能直接拿去写生产代码”的方向补齐:给你常用代码模板、易踩坑、JVM/打包/并发更深入的部分。(默认以 JDK 17 的习惯写;JDK 8 也基本通用,少量 API 需要替代。)
18. 生产级编码习惯(少踩坑=高产出)
18.1 null、Optional 与边界
- 对外接口(方法入参):能不收
null就不收,早校验早失败Objects.requireNonNull(x, "x must not be null")
- 对内流转:允许
null会让分支爆炸;更推荐:- 用默认值(如空集合
Collections.emptyList()) - 或用
Optional表达“可能没有”
- 用默认值(如空集合
- 注意:
Optional不适合作为字段/序列化载体(框架兼容性差),主要用于方法返回。
18.2 不可变对象(并发与可维护性神器)
- 能 immutable 就 immutable:字段
final、不暴露可变集合引用、构造时拷贝。 - 大数据里很常见的“诡异并发 bug”,很多是共享可变对象导致的。
18.3 record(如果你用 JDK 16+)
适合做 DTO/配置项/值对象(天然 equals/hashCode/toString):
public record TopicPartition(String topic, int partition) {}
但:如果需要可变字段、复杂校验、兼容某些序列化框架时要评估。
19. Stream API:该用就用,别用成性能黑洞
19.1 什么时候适合
- 数据量中等、逻辑清晰、链式变换可读性强
- 不在极致热点路径(例如每条日志解析都走 5 层 stream)
19.2 常见坑
stream().map(...).collect(...)产生很多中间对象- 千万慎用
parallelStream():- 默认 ForkJoinPool(可能和你业务线程池抢资源)
- 对 I/O、阻塞操作不友好
- 数据倾斜会更糟
19.3 替代建议
- 热点路径用 for 循环 + 预分配容器,通常更快更省内存。
20. 时间处理:只用 java.time
20.1 规则
- 存储/传输尽量用 UTC:
Instant、epoch millis - 展示给人看才转时区:
ZonedDateTime
Instant now = Instant.now(); long ts = now.toEpochMilli(); ZonedDateTime shanghai = now.atZone(ZoneId.of("Asia/Shanghai"));
20.2 常见坑
System.currentTimeMillis()可能回拨(NTP),做耗时用System.nanoTime()
21. 并发进阶(线程池、队列、原子类的正确姿势)
21.1 线程池:生产可用模板(强烈建议照着抄)
ThreadFactory tf = r -> { Thread t = new Thread(r); t.setName("etl-worker-" + t.threadId()); t.setDaemon(false); return t; }; BlockingQueue<Runnable> queue = new ArrayBlockingQueue<>(10_000); ThreadPoolExecutor pool = new ThreadPoolExecutor( 16, // core 32, // max 60, TimeUnit.SECONDS, queue, tf, new ThreadPoolExecutor.CallerRunsPolicy() // 背压:队列满就让提交方跑 ); // 建议:优雅停机时用 pool.shutdown();
要点:
- 队列必须有界(否则内存迟早被堆爆)
- 线程必须命名(排查时你会感谢自己)
- 拒绝策略要“有业务意义”:
CallerRunsPolicy常用于实现背压
21.2 原子类:计数优先 LongAdder
高并发统计计数:
AtomicLong:竞争激烈时会抖LongAdder:更适合高并发热点计数(metrics 常用)
java
Download
Copy code
LongAdder counter = new LongAdder(); counter.increment(); long v = counter.sum();
21.3 CompletableFuture:异步编排 + 超时兜底
ExecutorService ioPool = Executors.newFixedThreadPool(64); CompletableFuture<String> f = CompletableFuture.supplyAsync(() -> callRemote(), ioPool) .orTimeout(2, TimeUnit.SECONDS) .exceptionally(ex -> "fallback"); String result = f.join(); // join 不用强制捕获 checked exception
要点:一定要处理超时/异常,否则线上“卡死”等你来背锅。
21.4 共享变量可见性:volatile 的正确理解
volatile保证可见性与一定的有序性,但不保证复合操作原子性(如count++)- 计数用
AtomicLong/LongAdder;复合状态用锁或原子引用 + CAS 设计
22. 性能与内存:大数据 Java 的“省钱”技巧
22.1 少分配对象(降低 GC 压力)
- 热点路径避免:
- 频繁
new小对象 - 大量装箱(
Long/Integer) - 字符串拼接(尤其日志)
- 频繁
- 用原生类型流:
IntStream/LongStream(但也别滥用 stream)
22.2 集合容量预估(非常值钱)
new HashMap<>(expectedSize * 4 / 3 + 1)(负载因子默认 0.75)new ArrayList<>(expectedSize)
22.3 ByteBuffer/复用 buffer(I/O 场景常见)
- 读写网络/文件避免频繁分配 buffer
- 使用池化/复用(注意线程安全与生命周期)
22.4 “小心翼翼地”用第三方高性能集合
在极致性能场景,可考虑 fastutil/HPPC 这类 primitive 集合(减少装箱)。
但:引入依赖要评估与 Spark/Flink/Hadoop 的冲突风险(shade 可能要安排上)。
23. JVM 与 GC:你至少要能自救
23.1 常见 OOM 类型(看到就能定位方向)
java.lang.OutOfMemoryError: Java heap space:堆不够/泄漏/对象太多GC overhead limit exceeded:GC 忙到崩溃,通常也是堆压力OutOfMemoryError: Metaspace:类太多/重复加载/动态代理、依赖冲突OutOfMemoryError: Direct buffer memory:堆外内存(NIO/Netty)不足
23.2 GC 日志(JDK 11+ 推荐写法)
-XX:+UseG1GC -Xlog:gc*:file=gc.log:time,uptime,level,tags -XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=/data/dumps
你需要会看三件事:
- GC 频率是否异常(Young 太频繁/Full 出现)
- 暂停时间是否超 SLA(p99 pause)
- 老年代是否持续上涨(疑似泄漏或缓存无界)
23.3 线程问题定位(必会)
- CPU 飙高:
top -H找线程 →jstack对照 nid - 死锁:
jstack会直接提示 deadlock - 线程池排队:看队列长度、任务耗时、下游是否变慢
24. 日志与可观测性(让问题“自己开口说话”)
24.1 结构化日志 + MDC(链路定位神器)
- 在处理一条任务/一批数据时写入 trace 信息:
- jobId、batchId、topic、partition、offsetRange、requestId
- slf4j MDC(伪代码概念):
- 请求开始 set
- finally remove(线程池复用,必须清理)
24.2 指标建议(大数据链路通用)
- 输入:records/s、bytes/s
- 输出:sink 成功/失败数、重试数
- 延迟:处理耗时 histogram(p95/p99)
- 背压:队列长度、批次堆积量、consumer lag
- JVM:heap、GC pause、threads、direct memory
25. Maven 打包与依赖冲突:Spark/Flink 场景重点
25.1 provided 的典型用法
在 Spark/Flink 集群上运行,很多依赖集群已提供(例如某些 Hadoop 依赖)。你项目里用 provided 避免重复打包(减少冲突)——但要非常清楚“运行环境到底提供了什么”。
25.2 shade(胖包)常用模板(避免类冲突)
<plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <version>3.5.1</version> <executions> <execution> <phase>package</phase> <goals><goal>shade</goal></goals> <configuration> <createDependencyReducedPom>true</createDependencyReducedPom> <relocations> <relocation> <pattern>com.google.common</pattern> <shadedPattern>shade.yourapp.com.google.common</shadedPattern> </relocation> </relocations> </configuration> </execution> </executions> </plugin>
常见现象与结论:
- 运行时报
NoSuchMethodError:几乎必是依赖版本冲突 - 解决路径:
mvn dependency:tree→ 锁版本/排除依赖/relocation
26. 生产常见“玄学报错”对照表(其实都不玄学)
NoSuchMethodError:编译用的 A 版本,运行加载了 B 版本ClassNotFoundException:包没打进去 / scope 不对 / 类路径优先级问题ClassCastException:不同 classloader 加载了“同名不同类”(Spark/Flink 插件化时常见)- “本地好好的,上集群挂了”:十有八九是 依赖/类加载/配置 差异
27. 你可以直接照做的“训练营式任务”(带验收标准)
任务 A:有界队列的批处理管道(背压 + 指标)
目标:输入 1M 行文本,解析后批量写入(假写入也行),不 OOM、不无止境排队。
验收:
- 队列满时能背压(吞吐下降但不崩)
- 统计 p95/p99 处理耗时
- 可优雅停机(不丢正在处理的 batch)
任务 B:故障注入的重试策略(指数退避 + 抖动)
目标:模拟下游 10% 失败,重试最多 5 次,超限进入 DLQ(本地文件代替也行)。
验收:
- 重试间隔按指数增长(含随机抖动)
- 总失败数可观测(metrics + 日志字段)
任务 C:制造一次 OOM 并定位
目标:写一个无界缓存或装箱风暴 demo,触发 heap OOM,导出 heap dump,用工具找到“最大对象/最大引用链”。
验收:
- 能说明:是谁在持有这些对象、为什么无法回收、修复方案是什么
28. 一份“日常开发速查清单”(贴在工位上那种)
- 线程池:有界队列 + 命名 + 拒绝策略 + shutdown
- 外部调用:超时必须设、重试要限次、注意幂等
- 集合:预估容量;
HashMapkey 不可变 - 日志:参数化 + 上下文字段;线程池用 MDC 要记得 remove
- 序列化:别用 Java 原生序列化做跨系统协议
- 打包:先
dependency:tree再上集群;冲突用 shade/relocation
如果你还要“继续”,我可以在下一段把最实用的几块直接补成可复用代码组件(拿来即用的那种):
- 指数退避重试器(带抖动、最大时长、可取消)
- 批处理器(batching + flush + backpressure + 超时)
- 一个最小可运行的 Kafka consumer 模板(手动提交 offset + DLQ)
- 常用 JVM 启动参数模板(按:吞吐优先 / 延迟优先 两套)
我会直接给代码与注释,不需要你再提供信息。
下面给你 4 份“可直接复用到生产项目”的 Java 组件模板(偏大数据工程常用):重试器(指数退避+抖动)、批处理器(有界队列+背压+定时 flush)、Kafka Consumer(手动提交offset+再均衡处理+DLQ 思路)、JVM 启动参数模板(吞吐/低延迟两套)。
1) 指数退避重试器(Exponential Backoff + Jitter)
特点:
- 最大尝试次数、最大总耗时
- 指数退避 + 抖动(避免“惊群重试”)
- 可按异常类型/返回值决定是否重试
- 支持线程中断(不吞中断)
import java.time.Duration; import java.util.Objects; import java.util.concurrent.ThreadLocalRandom; import java.util.function.BiConsumer; import java.util.function.Predicate; public final class Retry { @FunctionalInterface public interface CheckedSupplier<T> { T get() throws Exception; } public static final class Policy { public final int maxAttempts; // 包含第一次调用 public final Duration baseDelay; // 第一次重试等待 public final Duration maxDelay; // 单次等待上限 public final Duration maxTotalTime; // 总耗时上限(含执行时间) public final double jitter; // 0.0~1.0 建议 0.2 public final Predicate<Exception> retryOn; // 哪些异常重试 public final BiConsumer<Integer, Exception> onRetry; // 第n次尝试失败回调(可打日志/指标) public Policy(int maxAttempts, Duration baseDelay, Duration maxDelay, Duration maxTotalTime, double jitter, Predicate<Exception> retryOn, BiConsumer<Integer, Exception> onRetry) { if (maxAttempts < 1) throw new IllegalArgumentException("maxAttempts must be >= 1"); this.maxAttempts = maxAttempts; this.baseDelay = Objects.requireNonNull(baseDelay); this.maxDelay = Objects.requireNonNull(maxDelay); this.maxTotalTime = Objects.requireNonNull(maxTotalTime); if (jitter < 0.0 || jitter > 1.0) throw new IllegalArgumentException("jitter must be in [0,1]"); this.jitter = jitter; this.retryOn = retryOn != null ? retryOn : (e -> true); this.onRetry = onRetry != null ? onRetry : (n, e) -> {}; } } public static <T> T call(CheckedSupplier<T> supplier, Policy policy) throws Exception { Objects.requireNonNull(supplier); Objects.requireNonNull(policy); final long startNanos = System.nanoTime(); Exception last = null; for (int attempt = 1; attempt <= policy.maxAttempts; attempt++) { try { return supplier.get(); } catch (Exception e) { last = e; boolean canRetry = attempt < policy.maxAttempts && policy.retryOn.test(e); policy.onRetry.accept(attempt, e); if (!canRetry) throw e; Duration sleep = backoffWithJitter(policy.baseDelay, policy.maxDelay, attempt, policy.jitter); long elapsedNanos = System.nanoTime() - startNanos; if (Duration.ofNanos(elapsedNanos).plus(sleep).compareTo(policy.maxTotalTime) > 0) { throw e; // 超过总耗时上限,直接失败 } sleepInterruptibly(sleep); } } throw last; // 理论上到不了 } private static Duration backoffWithJitter(Duration base, Duration max, int attempt, double jitter) { // attempt=1 表示第一次失败后的等待(也可以从 0 开始,这里简单点) long baseMs = Math.max(0, base.toMillis()); long maxMs = Math.max(0, max.toMillis()); // 指数:base * 2^(attempt-1) long expMs; if (attempt >= 31) expMs = Long.MAX_VALUE / 2; // 防溢出 else expMs = baseMs * (1L << (attempt - 1)); long capped = Math.min(expMs, maxMs); // 抖动:在 [capped*(1-j), capped*(1+j)] 范围随机 if (capped <= 0 || jitter <= 0.0) return Duration.ofMillis(Math.max(0, capped)); double factor = 1.0 + ThreadLocalRandom.current().nextDouble(-jitter, jitter); long jittered = (long) Math.max(0, capped * factor); return Duration.ofMillis(jittered); } private static void sleepInterruptibly(Duration d) { try { Thread.sleep(d.toMillis()); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); throw new RuntimeException("retry sleep interrupted", ie); } } }
用法示例(比如调用远程 HTTP/Kafka Admin/HDFS):
Retry.Policy p = new Retry.Policy( 5, Duration.ofMillis(200), Duration.ofSeconds(3), Duration.ofSeconds(10), 0.2, ex -> true, // 可改为:只对超时/5xx 重试 (attempt, ex) -> System.err.println("attempt=" + attempt + " failed: " + ex) ); String resp = Retry.call(() -> remoteCall(), p);
2) 批处理器(Batching + 有界队列背压 + 定时flush)
适用:日志解析后批量写 ES/ClickHouse/HTTP sink、批量写 DB、批量发 Kafka 等。
特点:
ArrayBlockingQueue有界:队列满时阻塞(天然背压)- 单独 worker 线程:按
maxBatchSize或maxWaitflush close()优雅停机:尽量 flush 剩余数据
import java.time.Duration; import java.util.ArrayList; import java.util.List; import java.util.Objects; import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.BlockingQueue; import java.util.concurrent.TimeUnit; public final class Batcher<T> implements AutoCloseable { public interface BatchSink<T> { void write(List<T> batch) throws Exception; } private final BlockingQueue<T> queue; private final int maxBatchSize; private final Duration maxWait; private final BatchSink<T> sink; private final Thread worker; private volatile boolean running = true; public Batcher(int queueCapacity, int maxBatchSize, Duration maxWait, BatchSink<T> sink, String workerName) { if (queueCapacity <= 0) throw new IllegalArgumentException("queueCapacity must be > 0"); if (maxBatchSize <= 0) throw new IllegalArgumentException("maxBatchSize must be > 0"); this.queue = new ArrayBlockingQueue<>(queueCapacity); this.maxBatchSize = maxBatchSize; this.maxWait = Objects.requireNonNull(maxWait); this.sink = Objects.requireNonNull(sink); this.worker = new Thread(this::runLoop); this.worker.setName(workerName != null ? workerName : "batcher-worker"); this.worker.setDaemon(false); this.worker.start(); } // 阻塞提交:队列满则背压(你也可以改成 offer+timeout+丢弃策略) public void submit(T item) throws InterruptedException { if (!running) throw new IllegalStateException("batcher already closed"); queue.put(item); } private void runLoop() { List<T> batch = new ArrayList<>(maxBatchSize); long deadlineNanos = System.nanoTime() + maxWait.toNanos(); while (running || !queue.isEmpty()) { try { long waitNanos = Math.max(0, deadlineNanos - System.nanoTime()); T item = queue.poll(waitNanos, TimeUnit.NANOSECONDS); if (item != null) { batch.add(item); } boolean sizeReached = batch.size() >= maxBatchSize; boolean timeReached = System.nanoTime() >= deadlineNanos; if (!batch.isEmpty() && (sizeReached || timeReached)) { flush(batch); batch.clear(); deadlineNanos = System.nanoTime() + maxWait.toNanos(); } else if (batch.isEmpty() && timeReached) { deadlineNanos = System.nanoTime() + maxWait.toNanos(); } } catch (InterruptedException ie) { Thread.currentThread().interrupt(); // 被中断时:尝试 flush 已有数据,然后退出 safeFlush(batch); return; } catch (Exception e) { // sink 写入失败:这里给一个“不中断主循环”的兜底 // 生产建议:记录日志 + 指标 + 必要时将 batch 写入本地 DLQ 文件 safeFlush(batch); // 你也可以选择不 flush 或者重试 batch.clear(); deadlineNanos = System.nanoTime() + maxWait.toNanos(); } } safeFlush(batch); } private void flush(List<T> batch) throws Exception { sink.write(batch); } private void safeFlush(List<T> batch) { if (batch == null || batch.isEmpty()) return; try { sink.write(batch); } catch (Exception ignored) { // 最后兜底:避免 close 时抛异常导致资源未释放 } } @Override public void close() { running = false; worker.interrupt(); try { worker.join(); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); } } }
用法示例(批量写入):
Batcher<String> b = new Batcher<>( 10_000, 500, Duration.ofSeconds(1), batch -> bulkWrite(batch), "sink-batcher" ); b.submit("row1"); // ... b.close();
3) Kafka Consumer 模板(手动提交 offset + rebalance 处理 + DLQ 思路)
目标:最常见的“稳定消费”骨架。
enable.auto.commit=false- 每次处理成功后
commitSync(或批量提交) - 处理 rebalance:在 partitions revoke 时提交已处理 offset
- 失败数据进 DLQ(演示思路;你可以改成本地文件/另一个 topic)
依赖:
org.apache.kafka:kafka-clients
import org.apache.kafka.clients.consumer.*; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.WakeupException; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; import java.time.Duration; import java.util.*; import java.util.concurrent.atomic.AtomicBoolean; public final class KafkaConsumerApp { public static void main(String[] args) { String bootstrap = "localhost:9092"; String topic = "input-topic"; String groupId = "etl-group"; String dlqTopic = "input-topic-dlq"; Properties cp = new Properties(); cp.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrap); cp.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); cp.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); cp.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); cp.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); cp.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "500"); cp.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); // cp.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, ""); // 处理很慢时需要调整,避免被踢出组 // cp.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, ""); // 与心跳相关,按需调整 Properties pp = new Properties(); pp.put("bootstrap.servers", bootstrap); pp.put("key.serializer", StringSerializer.class.getName()); pp.put("value.serializer", StringSerializer.class.getName()); // 生产建议:acks、retries、idempotence 等按可靠性要求配置 AtomicBoolean running = new AtomicBoolean(true); try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(cp); KafkaProducer<String, String> dlqProducer = new KafkaProducer<>(pp)) { Map<TopicPartition, OffsetAndMetadata> offsetsToCommit = new HashMap<>(); ConsumerRebalanceListener rebalanceListener = new ConsumerRebalanceListener() { @Override public void onPartitionsRevoked(Collection<TopicPartition> partitions) { // rebalance 前提交已处理 offset,尽量减少重复消费 if (!offsetsToCommit.isEmpty()) { consumer.commitSync(offsetsToCommit); } } @Override public void onPartitionsAssigned(Collection<TopicPartition> partitions) { // 可在此做 seek、恢复本地状态等 } }; consumer.subscribe(Collections.singletonList(topic), rebalanceListener); Runtime.getRuntime().addShutdownHook(new Thread(() -> { running.set(false); consumer.wakeup(); // 让 poll 立刻退出 })); while (running.get()) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(1)); for (ConsumerRecord<String, String> r : records) { try { // 你的业务处理:解析、校验、写 sink ... process(r.key(), r.value()); // 记录“下一个要提交的 offset”(提交的是 next offset) offsetsToCommit.put( new TopicPartition(r.topic(), r.partition()), new OffsetAndMetadata(r.offset() + 1) ); } catch (Exception ex) { // 失败处理:写 DLQ(或本地文件),避免卡住分区 String dlqValue = "error=" + ex.getClass().getSimpleName() + ", msg=" + ex.getMessage() + ", orig=" + r.value(); dlqProducer.send(new ProducerRecord<>(dlqTopic, r.key(), dlqValue)); // 根据语义决定:失败了是否提交 offset? // - 想“跳过坏数据”就提交 offset(避免永远重试同一条) // - 想“必须成功”就别提交,让它重试(但要配合重试上限+告警,否则会卡死) offsetsToCommit.put( new TopicPartition(r.topic(), r.partition()), new OffsetAndMetadata(r.offset() + 1) ); } } if (!offsetsToCommit.isEmpty()) { consumer.commitSync(offsetsToCommit); offsetsToCommit.clear(); } } } catch (WakeupException we) { // 正常退出路径 } } private static void process(String key, String value) { // TODO your logic if (value == null) throw new IllegalArgumentException("null value"); } }
你后续做“更稳更快”的版本时,通常会加:
- 批处理 + sink bulk(配合上面 Batcher)
commitAsync(吞吐更高)+ 周期性commitSync保底- 更细的重试策略(用上面 Retry)
- 处理耗时很长时:必须关注
max.poll.interval.ms,否则 rebalance 频繁
4) JVM 启动参数模板(大数据任务常用两套)
4.1 通用必配(排障友好)
-XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=/data/dumps -Xlog:gc*:file=/data/logs/gc.log:time,uptime,level,tags -Dfile.encoding=UTF-8
4.2 吞吐优先(批处理/离线 ETL 常见,G1)
-XX:+UseG1GC -Xms4g -Xmx4g -XX:MaxGCPauseMillis=200 -XX:+ParallelRefProcEnabled
4.3 低延迟优先(实时链路/服务化组件更常见,G1 调参偏“稳”)
-XX:+UseG1GC -Xms2g -Xmx2g -XX:MaxGCPauseMillis=50 -XX:G1NewSizePercent=30 -XX:G1MaxNewSizePercent=60 -XX:+ParallelRefProcEnabled
4.4 容器环境(K8s/YARN Container)建议
如果你不想写死 -Xmx,可以用百分比(JDK 10+):
-XX:MaxRAMPercentage=75 -XX:InitialRAMPercentage=75
但大数据任务里很多团队仍倾向显式 -Xms/-Xmx,更可控、也更方便复现实验。
更多推荐
所有评论(0)