如何在Scala开发的Akka-http中通过JDBC操作MySQL数据库
Scala Akka HTTP项目JDBC对接MySQL:原生SQL执行与库表创建实现
前置说明:以下实现默认你已经拿到了可用的
java.sql.Connection实例(不管是原生DriverManager直连还是HikariCP等连接池托管的连接都通用)。
一、数据库与数据表创建
建库建表属于DDL操作,JDBC默认配置下执行后会自动提交,不需要额外手动提交事务,注意执行建库语句时的连接不要提前指定不存在的库名,否则会报库不存在的错误。
- 执行建库
直接通过Statement执行CREATE DATABASE语句即可,建议显式指定字符集为utf8mb4,避免emoji和特殊字符存储乱码:
import java.sql.Connection val conn: Connection = ??? // 你已经实现好的连接获取逻辑 val stmt = conn.createStatement() try { val createDbSql = """ |CREATE DATABASE IF NOT EXISTS biz_default |DEFAULT CHARACTER SET utf8mb4 |COLLATE utf8mb4_unicode_ci |""".stripMargin stmt.executeUpdate(createDbSql) } finally { stmt.close() // 确保Statement资源释放 }
- 执行建表
建库完成后,可以选择两种方式指定操作的目标库:一是重新获取JDBC连接时在url中拼接对应参数指定库,二是执行USE biz_default语句切换当前连接的操作库,之后再执行建表语句:
val stmtForTable = conn.createStatement() try { stmtForTable.executeUpdate("USE biz_default") // 连接url已指定库可跳过 val createTableSql = """ |CREATE TABLE IF NOT EXISTS user_info ( | id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT PRIMARY KEY, | username VARCHAR(32) NOT NULL UNIQUE COMMENT '用户名', | age INT UNSIGNED NOT NULL DEFAULT 0 COMMENT '年龄', | create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间' |) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='用户基础信息表' |""".stripMargin stmtForTable.executeUpdate(createTableSql) } finally { stmtForTable.close() }
注意:生产环境不要在服务每次启动时无脑重复执行建库建表语句,建议搭配版本化DDL工具做表结构变更,或者仅在服务首次初始化时执行DDL,避免高权限账号带来的误操作风险。
二、原生SQL执行方法
所有JDBC操作都是阻塞IO,在Akka HTTP体系下不要把这类操作直接提交到默认的Actor调度线程池,需要单独配置专属的阻塞IO线程池来承载JDBC逻辑,避免占满核心线程导致服务吞吐暴跌。
- 增删改类(INSERT/UPDATE/DELETE)操作
优先使用PreparedStatement做参数占位,不要直接拼接SQL参数值,从根源避免SQL注入风险,注意JDBC的参数下标从1开始计数,不是从0开始:
val insertSql = "INSERT INTO user_info (username, age) VALUES (?, ?)" val pstmt = conn.prepareStatement(insertSql) try { pstmt.setString(1, "lisi") pstmt.setInt(2, 28) val affectedRows = pstmt.executeUpdate() // 返回值为SQL执行影响的行数 println(s"插入操作完成,受影响行数:$affectedRows") } finally { pstmt.close() }
- 查询类(SELECT)操作
执行查询后拿到ResultSet结果集,逐行遍历映射为Scala对象即可,后续可以直接传入Akka HTTP的路由层做JSON序列化返回:
// 先定义结果映射的样例类 case class User( id: Long, username: String, age: Int, createTime: java.sql.Timestamp ) val querySql = "SELECT id, username, age, create_time FROM user_info WHERE age > ?" val queryPstmt = conn.prepareStatement(querySql) try { queryPstmt.setInt(1, 18) val rs = queryPstmt.executeQuery() var result = List.empty[User] while (rs.next()) { val user = User( id = rs.getLong("id"), username = rs.getString("username"), age = rs.getInt("age"), createTime = rs.getTimestamp("create_time") ) result = user :: result } // result即为查询到的所有符合条件的用户数据 } finally { queryPstmt.close() conn.close() // 如果是连接池模式,这里的close是把连接归还到池里,不是物理断开 }
内容的提问来源于stack exchange,提问作者Anil Singh
相关产品推荐
相关产品推荐

