Kafka与Spark Streaming处理JSON:能否用foreachRDD提取字段写入MySQL?
Time: 1527776685000 ms
{"logtyp":"ERROR","LogTypName":"app.warning.exception","LogZeitpunkt":"Thu May 31 16:24:42 CEST 2018"}
{"logtyp":"ERROR","LogTypName":"app.warning.exception","LogZeitpunkt":"Thu May 31 16:24:44 CEST 2018"}
Time: 1527776690000 ms
{"logtyp":"ERROR","LogTypName":"app.warning.exception","LogZeitpunkt":"Thu May 31 16:24:45 CEST 2018"}
{"logtyp":"ERROR","LogTypName":"app.warning.exception","LogZeitpunkt":"Thu May 31 16:24:46 CEST 2018"}
现咨询是否可以使用`foreachRDD`方法提取JSON中的字段,从而将每条JSON记录以如下SQL语句形式写入MySQL: ```sql insert into my_table (logtyp, logtypname, logzeitpunkt) values ("ERROR", "app.warning.exception", "Thu May 31 16:24:46 CEST 2018");
解决方案
当然可以!这是Spark Streaming处理Kafka数据后落地关系型数据库的典型场景,我来一步步给你实现:
1. 准备JSON解析工具
首先需要能解析JSON字符串,这里推荐用Jackson(Spark内置依赖,不用额外引入太多包),我们先定义一个简单的实体类来映射JSON字段,这样代码更清晰,避免硬编码字段名:
import com.fasterxml.jackson.annotation.JsonProperty; public class LogRecord { @JsonProperty("logtyp") private String logtyp; @JsonProperty("LogTypName") private String logTypName; @JsonProperty("LogZeitpunkt") private String logZeitpunkt; // 生成getter和setter方法 public String getLogtyp() { return logtyp; } public void setLogtyp(String logtyp) { this.logtyp = logtyp; } public String getLogTypName() { return logTypName; } public void setLogTypName(String logTypName) { this.logTypName = logTypName; } public String getLogZeitpunkt() { return logZeitpunkt; } public void setLogZeitpunkt(String logZeitpunkt) { this.logZeitpunkt = logZeitpunkt; } }
2. 使用foreachRDD+foreachPartition写入MySQL
这里有个关键注意点:绝对不要在foreach里直接创建数据库连接,会导致每条数据都创建一个连接,性能极差甚至把数据库打垮。正确的做法是在foreachPartition里创建一个连接,整个分区共用这个连接,大大减少连接开销。
完整代码示例如下:
import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.spark.api.java.function.VoidFunction; import java.sql.Connection; import java.sql.DriverManager; import java.sql.PreparedStatement; import java.util.Iterator; // 初始化Jackson的ObjectMapper(可以在Driver端初始化,序列化后传到Executor) ObjectMapper objectMapper = new ObjectMapper(); jsonline.foreachRDD(new VoidFunction<JavaRDD<String>>() { @Override public void call(JavaRDD<String> rdd) throws Exception { // 跳过空RDD,避免无数据时创建连接 if (!rdd.isEmpty()) { rdd.foreachPartition(new VoidFunction<Iterator<String>>() { @Override public void call(Iterator<String> jsonIterator) throws Exception { // 1. 在Executor端创建数据库连接(每个分区一个连接) Connection conn = null; PreparedStatement stmt = null; try { // 加载MySQL驱动 Class.forName("com.mysql.cj.jdbc.Driver"); // 建立连接(建议把配置抽成常量或配置文件) conn = DriverManager.getConnection( "jdbc:mysql://your-mysql-host:3306/your-db", "username", "password" ); // 2. 预编译SQL,避免SQL注入,提升性能 String sql = "INSERT INTO my_table (logtyp, logtypname, logzeitpunkt) VALUES (?, ?, ?)"; stmt = conn.prepareStatement(sql); // 3. 遍历分区内的每条JSON数据 while (jsonIterator.hasNext()) { String jsonStr = jsonIterator.next(); try { // 解析JSON为实体类 LogRecord log = objectMapper.readValue(jsonStr, LogRecord.class); // 设置SQL参数 stmt.setString(1, log.getLogtyp()); stmt.setString(2, log.getLogTypName()); stmt.setString(3, log.getLogZeitpunkt()); // 添加到批处理(可选,批量插入性能更好) stmt.addBatch(); } catch (Exception e) { // 处理JSON解析错误,比如打印日志跳过错误数据 System.err.println("解析JSON失败: " + jsonStr + ", 错误信息: " + e.getMessage()); } } // 执行批处理插入 stmt.executeBatch(); } finally { // 4. 关闭资源,避免连接泄漏 if (stmt != null) stmt.close(); if (conn != null) conn.close(); } } }); } } });
3. 额外优化建议
- 使用连接池:上面的示例是直接创建连接,生产环境建议用连接池(比如HikariCP),避免频繁创建销毁连接的开销。
- 异常处理:可以把错误数据收集起来,比如写入HDFS或其他存储,方便后续排查。
- 事务控制:如果需要保证分区内数据的一致性,可以在打开连接后开启事务,执行完批处理后提交,失败则回滚。
- 参数配置:把MySQL的连接信息、批处理大小等参数放到配置文件里,不要硬编码在代码中。
内容的提问来源于stack exchange,提问作者Alex
相关产品推荐
相关产品推荐

