下面是一份“面向大数据开发工程师”的 Java 学习笔记(更偏工程落地 + 性能 + JVM + 并发 + 生态集成),学完后基本能覆盖你在 Hadoop/Spark/Flink/Kafka/ES/OLAP 等场景里写生产级 Java 的核心需求。你可以把它当作:路线图 + 速查笔记 + 易错点清单 + 练习清单


0. 学习目标(大数据工程视角)

你需要的 Java 能力通常不是“会语法就行”,而是:

  1. 写得对:类型系统、集合、序列化、并发安全、资源释放、异常边界清晰
  2. 跑得稳:JVM 内存/GC/类加载、线程模型、反压与超时、日志与指标
  3. 跑得快:对象分配、装箱拆箱、集合选择、I/O 模型、锁竞争、热点优化
  4. 好交付:Maven/Gradle、依赖冲突、shade、可测试性、可观测性
  5. 能对接生态: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 一致性
    • hashCodeequals 兼容
  • 常见坑:对象字段可变导致 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)
  • 你需要至少理解:
    • ByteBufferposition/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 看冲突
  • scopecompile/provided/test/runtime
  • shade/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. 一份“学完就能上手”的练习清单(建议按顺序做)

  1. 实现一个高性能日志解析器:支持 gzip 输入、逐行处理、输出统计指标
  2. 写一个 Kafka Consumer 服务:手动提交 offset + 重试 + 死信队列(DLQ)
  3. 实现一个并发批处理器:线程池 + 有界队列 + backpressure + 指标
  4. 做一次 JVM 排障演练:故意制造内存泄漏/线程死锁,然后用工具定位
  5. 打包一个可在 Spark/Flink 集群跑的 fat-jar:解决一次依赖冲突(shade relocation)
  6. 实现 Avro/Protobuf schema 演进 demo:验证向后/向前兼容

17. 面试/自检速查(你可以当 checklist)

  • HashMap 原理、扩容、为什么 key 不能乱变
  • 线程池 7 参数、为什么无界队列危险
  • volatile 能解决什么不能解决什么
  • JVM:堆/元空间/直接内存,OOM 类型与排查思路
  • GC 日志能看懂哪些关键信号
  • 序列化方案对比:JSON vs Avro/Protobuf vs Parquet
  • Maven 依赖冲突定位与解决
  • 日志与 metrics 如何设计(你会比只会写业务的人强很多)

推荐学习顺序(最省时间的路径)

  1. 集合/泛型/异常/资源管理(打底)
  2. 并发 + 线程池(决定上限)
  3. JVM + GC + 工具(决定你能不能救火)
  4. I/O + 序列化 + 构建依赖(决定吞吐与交付)
  5. 框架对接与工程化(Spark/Flink/Kafka 实战)

下面继续把这份笔记往“能直接拿去写生产代码”的方向补齐:给你常用代码模板、易踩坑、JVM/打包/并发更深入的部分。(默认以 JDK 17 的习惯写;JDK 8 也基本通用,少量 API 需要替代。)


18. 生产级编码习惯(少踩坑=高产出)

18.1 nullOptional 与边界

  • 对外接口(方法入参):能不收 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 规则

  • 存储/传输尽量用 UTCInstant、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

你需要会看三件事:

  1. GC 频率是否异常(Young 太频繁/Full 出现)
  2. 暂停时间是否超 SLA(p99 pause)
  3. 老年代是否持续上涨(疑似泄漏或缓存无界)

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
  • 外部调用:超时必须设、重试要限次、注意幂等
  • 集合:预估容量;HashMap key 不可变
  • 日志:参数化 + 上下文字段;线程池用 MDC 要记得 remove
  • 序列化:别用 Java 原生序列化做跨系统协议
  • 打包:先 dependency:tree 再上集群;冲突用 shade/relocation

如果你还要“继续”,我可以在下一段把最实用的几块直接补成可复用代码组件(拿来即用的那种):

  1. 指数退避重试器(带抖动、最大时长、可取消)
  2. 批处理器(batching + flush + backpressure + 超时)
  3. 一个最小可运行的 Kafka consumer 模板(手动提交 offset + DLQ)
  4. 常用 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 线程:按 maxBatchSizemaxWait flush
  • 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,更可控、也更方便复现实验。


更多推荐