Erlang 应用访问关系型数据库的方式与其他语言有本质差异:没有 ORM 的对象映射负担,没有线程本地连接的概念,连接是进程持有的资源。这带来两个后果——连接池必须显式管理(否则进程数一涨就会打爆数据库),事务必须绑定在单个进程与单条连接上(否则跨进程共享连接会串事务)。理解这两点,是用好 Erlang 数据库集成的关键。本文将以 PostgreSQL 的 epgsql 与 MySQL 的 mysql-otp 为主线,讲解驱动选型、poolboy 连接池、事务与预处理语句、迁移与索引策略,以及超时与故障重连的治理。
一、Erlang 数据库访问生态
1.1 生态概览
| 层次 | 方案 | 说明 |
|---|---|---|
| 原生驱动 | epgsql、mysql-otp | 纯 Erlang,直接讲协议 |
| 查询构建 | pgo、sqerl | 提供池化与查询封装 |
| Elixir ORM | Ecto | Elixir 侧的事实标准 |
| 连接池 | poolboy、pgo 内置池 | 独立进程池 |
| 迁移工具 | dbmate、Ecto Migration | 独立于应用运行 |
重要前提:无论用哪种驱动,连接都是进程绑定的资源。一个进程在同一时刻只能安全使用一条连接,跨进程共享连接会导致事务串台与协议错乱。
1.2 连接模型对比
| 维度 | Erlang 原生驱动 | 典型 ORM 语言 |
|---|---|---|
| 连接持有者 | 进程 | 线程 |
| 池化方式 | 显式进程池 | 框架内置 |
| 事务范围 | 单连接,手动传递 | 隐式上下文 |
| 结果映射 | tuple / map | 对象 |
| 超时控制 | 每调用显式指定 | 全局配置 |
二、驱动选型:epgsql 与 mysql-otp
2.1 epgsql 连接与查询
%% rebar.config
{deps, [{epgsql, "4.7.1"}, {poolboy, "1.5.2"}]}.
-module(pg_demo).
-export([connect/0, query/2]).
connect() ->
{ok, Conn} = epgsql:connect(#{
host => "10.0.1.5",
port => 5432,
username => "app",
password => "secret",
database => "prod",
timeout => 5000,
ssl => true
}),
Conn.
%% 参数化查询:$1 / $2 占位,自动转义,杜绝注入
find_user(Conn, UserId) ->
case epgsql:equery(Conn, "SELECT id, name, email FROM users WHERE id = $1",
[UserId]) of
{ok, _Cols, [{Id, Name, Email}]} ->
{ok, #{id => Id, name => Name, email => Email}};
{ok, _Cols, []} ->
{error, not_found};
{error, Reason} ->
{error, Reason}
end.
2.2 mysql-otp 连接与查询
{ok, Pid} = mysql:start_link([
{host, "10.0.1.6"},
{port, 3306},
{user, "app"},
{password, "secret"},
{database, "prod"},
{query_timeout, 5000},
{connect_timeout, 5000}
]).
%% 参数化查询:? 占位
mysql:query(Pid, "SELECT id, name FROM users WHERE id = ?", [UserId]).
2.3 驱动能力对比
| 能力 | epgsql | mysql-otp |
|---|---|---|
| 占位符 | $1 $2 | ? |
| 返回格式 | tuple / map({column, ...}) | 列表 of maps |
| 预处理语句 | 支持(parse/prepared_query) | 支持(prepare) |
| 通知/监听 | LISTEN/NOTIFY | 无 |
| 大对象 | COPY 支持 | LOAD DATA |
| SSL | 支持 | 支持 |
选型建议:PostgreSQL 优先 epgsql(支持
LISTEN/NOTIFY、COPY、丰富类型);MySQL 用 mysql-otp。若使用 Elixir,直接用 Ecto(见 https://plumephp.com/elixir-ecto-data-access/)而非裸驱动。
2.4 类型映射陷阱
%% PostgreSQL 的 numeric 默认返回 binary,需显式转换
{ok, _, [{<<"123.45">>}]} = epgsql:equery(Conn, "SELECT 123.45::numeric"),
%% 解决方案:查询时转换,或在驱动层配置
epgsql:equery(Conn, "SELECT amount::float8 FROM orders WHERE id = $1", [Id]).
%% timestamp 返回 {{Y,M,D},{H,Mi,S}} 元组
{ok, _, [{{2026,10,1},{12,0,0}}]} = epgsql:equery(Conn, "SELECT now()::timestamp").
三、连接池设计
3.1 为什么必须用连接池
每个数据库连接占用服务端一个后端进程(PostgreSQL 每连接一个进程,内存开销约 5~10 MB)。若每个 Erlang 请求进程都新建连接,几百并发就能拖垮数据库:
%% 反模式:每次请求新建连接
handle_request(Req) ->
{ok, Conn} = epgsql:connect(Opts), % 昂贵!
Result = epgsql:equery(Conn, Sql, Params),
epgsql:close(Conn),
Result.
%% 正解:从池中借出,用完归还
handle_request(Req) ->
poolboy:transaction(db_pool, fun(Worker) ->
gen_server:call(Worker, {query, Sql, Params})
end).
3.2 poolboy 池的监督树接入
%% 池规格(放入应用监督树)
PoolArgs = [
{name, {local, db_pool}},
{worker_module, pg_worker},
{size, 16}, % 常驻连接数
{max_overflow, 8} % 峰值可临时新建
],
ChildSpec = poolboy:child_spec(db_pool, PoolArgs, PgConnectArgs),
{ok, {SupFlags, [ChildSpec]}}.
%% pg_worker.erl:每 worker 持有一条连接
-module(pg_worker).
-behaviour(gen_server).
-behaviour(poolboy_worker).
-export([init/1, handle_call/3, handle_cast/2, terminate/2]).
init(ConnectArgs) ->
{ok, Conn} = epgsql:connect(ConnectArgs),
{ok, #{conn => Conn}}.
handle_call({query, Sql, Params}, _From, #{conn := Conn} = State) ->
{reply, epgsql:equery(Conn, Sql, Params), State};
handle_call(_Req, _From, State) ->
{reply, {error, unknown}, State}.
handle_cast(_Msg, State) -> {noreply, State}.
terminate(_Reason, #{conn := Conn}) -> catch epgsql:close(Conn), ok.
3.3 池容量规划
| 参数 | 含义 | 建议 |
|---|---|---|
size | 常驻连接数 | 约 核数 × 2,不超过 DB max_connections 的 1/N |
max_overflow | 峰值溢出 | size 的 25%~50% |
| 排队超时 | 借出等待上限 | 1~5 秒,超时快速失败 |
| 单查询超时 | 语句执行上限 | 按业务 SLA,避免长查询占连接 |
%% 借出时设置等待超时,避免请求无限挂起
case poolboy:checkout(db_pool, false, 1000) of
Worker when is_pid(Worker) ->
try gen_server:call(Worker, {query, Sql, Params})
after poolboy:checkin(db_pool, Worker)
end;
full ->
{error, pool_exhausted}
end.
容量陷阱:多个 Erlang 节点共享一个数据库时,总连接数 = 节点数 × 池大小。务必用
节点数 × (size + max_overflow) < DB max_connections × 0.8校验,否则扩容节点会打爆数据库。
3.4 使用 ETS 缓存降低池压力
高频只读查询可先走 ETS 缓存(见 https://plumephp.com/erlang-ets-caching/),未命中才占用连接:
get_user_cached(UserId) ->
case ets:lookup(user_cache, UserId) of
[{_, User}] -> {ok, User};
[] ->
{ok, User} = poolboy:transaction(db_pool, fun(W) ->
gen_server:call(W, {query_user, UserId})
end),
ets:insert(user_cache, {UserId, User}),
{ok, User}
end.
四、事务与预处理语句
4.1 显式事务
%% 事务必须绑定在同一条连接上,因此整个事务在 worker 进程内完成
transfer(From, To, Amount) ->
poolboy:transaction(db_pool, fun(Worker) ->
gen_server:call(Worker, {transaction, fun(Conn) ->
{ok, _, _} = epgsql:equery(Conn,
"UPDATE accounts SET balance = balance - $1 WHERE id = $2",
[Amount, From]),
{ok, _, _} = epgsql:equery(Conn,
"UPDATE accounts SET balance = balance + $1 WHERE id = $2",
[Amount, To]),
ok
end}, 10000)
end).
%% worker 侧:用 sqerl 风格的事务包装(伪代码示意语义)
handle_call({transaction, Fun}, _From, #{conn := Conn} = State) ->
Result = epgsql:with_transaction(Conn, fun(C) -> Fun(C) end, [{reraise, true}]),
{reply, Result, State}.
4.2 事务的常见陷阱
| 陷阱 | 后果 | 规避 |
|---|---|---|
| 事务中调用外部服务 | 连接被长时间占用 | 事务内只做 DB 操作 |
| 跨进程共享连接 | 事务串台 | 事务只在 worker 内完成 |
| 忘记回滚 | 连接残留脏状态 | 用 with_transaction 自动回滚 |
| 长事务持有锁 | 阻塞其他写入 | 控制事务粒度,设置 statement_timeout |
%% PostgreSQL:服务端强制语句超时,防止长查询拖垮池
%% SET statement_timeout = '5s';
%% SET idle_in_transaction_session_timeout = '10s';
4.3 预处理语句
预处理语句避免重复解析 SQL,并从根本上杜绝注入:
%% 1. 解析(返回语句名与参数类型)
{ok, Statement} = epgsql:parse(Conn,
"INSERT INTO events (type, payload) VALUES ($1, $2)", []),
%% 2. 多次执行(复用已解析的计划)
epgsql:prepared_query(Conn, Statement, [<<"order.created">>, Payload1]),
epgsql:prepared_query(Conn, Statement, [<<"order.paid">>, Payload2]),
%% 3. 批量插入(COPY 比逐条 insert 快一个数量级)
epgsql:copy_from(Conn, "COPY events (type, payload) FROM STDIN",
[{<<"order.created">>, <<"{}">>} | MoreRows]).
4.4 批处理与分页
%% 批量插入:拼一条多值 INSERT,比 N 次单条快得多
batch_insert(Conn, Rows) ->
Placeholders = string:join(
[io_lib:format("($~p, $~p)", [I * 2 + 1, I * 2 + 2])
|| I <- lists:seq(0, length(Rows) - 1)], ", "),
Sql = "INSERT INTO events (type, payload) VALUES " ++ Placeholders,
Params = lists:flatten(Rows),
epgsql:equery(Conn, Sql, Params).
%% 大结果集分页,避免一次性加载
page(Conn, Limit, Offset) ->
epgsql:equery(Conn,
"SELECT id, name FROM users ORDER BY id LIMIT $1 OFFSET $2",
[Limit, Offset]).
五、迁移、索引与故障治理
5.1 迁移策略
Erlang 生态没有统一迁移框架,推荐用独立工具并把迁移脚本纳入版本控制:
-- migrations/20261001_create_orders.sql
CREATE TABLE IF NOT EXISTS orders (
id BIGSERIAL PRIMARY KEY,
user_id BIGINT NOT NULL REFERENCES users(id),
amount NUMERIC(12,2) NOT NULL,
status TEXT NOT NULL DEFAULT 'pending',
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_orders_user_created
ON orders (user_id, created_at DESC);
# dbmate 常用命令
dbmate new create_orders # 生成迁移文件
dbmate up # 执行未应用的迁移
dbmate status # 查看状态
dbmate rollback # 回滚最近一次
迁移原则:迁移必须是向前兼容的(应用新旧版本能同时跑),因此采用「扩展-迁移-收缩」三步法:先加新列(可空)、再双写、最后删旧列。
5.2 索引设计要点
| 场景 | 索引类型 | 说明 |
|---|---|---|
| 等值查询 | B-tree | 默认,最通用 |
| 范围 + 排序 | 复合 B-tree | 列顺序按「等值列在前、范围列在后」 |
| 模糊匹配 | GIN + trigram | pg_trgm 扩展 |
| JSONB 字段 | GIN | jsonb_path_ops 更省空间 |
| 唯一约束 | UNIQUE | 兼作幂等去重手段 |
%% 用 EXPLAIN 验证索引是否生效
{ok, _, Rows} = epgsql:equery(Conn,
"EXPLAIN (ANALYZE, BUFFERS) SELECT * FROM orders WHERE user_id = $1", [UserId]),
%% 关注输出中的 Index Scan / Seq Scan / Rows Removed by Filter
5.3 超时治理
超时必须在三个层次同时设置,缺一层就会出现「连接被永久占用」:
%% 1. 连接建立超时
{ok, Conn} = epgsql:connect(#{timeout => 5000, ...}),
%% 2. 语句执行超时(客户端侧)
epgsql:equery(Conn, Sql, Params, 5000),
%% 3. 服务端侧超时(最可靠,连接归还后自动清理)
%% SET statement_timeout = '5s';
| 层次 | 参数 | 建议值 |
|---|---|---|
| 建连 | connect_timeout | 5s |
| 查询 | query_timeout | 按 SLA,通常 3~5s |
| 借出等待 | checkout 超时 | 1s |
| 服务端 | statement_timeout | 略大于客户端 |
5.4 故障重连与降级
%% 连接断开时 worker 崩溃,由 poolboy 监督树自动重建
%% 关键:确保 terminate 里关闭旧连接,init 里新建连接
init(ConnectArgs) ->
case epgsql:connect(ConnectArgs) of
{ok, Conn} -> {ok, #{conn => Conn}};
{error, Reason} -> {stop, {connect_failed, Reason}}
end.
terminate(_Reason, #{conn := Conn}) ->
catch epgsql:close(Conn), % 幂等关闭,忽略已断开错误
ok.
%% 数据库不可用时的降级:池满或超时都快速失败,不无限等待
get_user(UserId) ->
case poolboy:checkout(db_pool, false, 500) of
full -> {error, service_unavailable};
Worker ->
try gen_server:call(Worker, {query_user, UserId}, 3000)
catch exit:{timeout, _} -> {error, timeout}
after poolboy:checkin(db_pool, Worker)
end
end.
熔断思路:当连续失败超过阈值时,短暂拒绝所有请求(快速失败),给数据库恢复窗口。这比让所有请求排队等待超时更有利于系统整体恢复。可结合 https://plumephp.com/erlang-otp-framework/ 中的监督树,把「健康探测」做成独立进程。
六、总结
Erlang 的数据库集成没有魔法,核心约束只有两条:连接是进程绑定的,事务必须绑定单连接。围绕这两条约束,工程实践可以归纳为:
- 驱动选型:PostgreSQL 用 epgsql,MySQL 用 mysql-otp,Elixir 项目直接用 Ecto;始终用参数化占位符,杜绝拼接 SQL;
- 连接池:用 poolboy 把连接数控制在数据库可承受范围内,容量规划要算上所有节点;
- 事务:事务只在单个 worker 进程内完成,用
with_transaction保证回滚,事务内绝不调用外部服务; - 预处理与批处理:复用已解析语句,批量插入用多值 INSERT 或
COPY; - 迁移与索引:迁移向前兼容(扩展-迁移-收缩),索引按查询模式设计,用
EXPLAIN验证; - 超时与降级:建连、查询、借出等待、服务端四层超时全设,配合熔断与快速失败。
把连接池当作「有界的稀缺资源」来管理,把事务当作「不可跨越进程边界的原子块」来使用,Erlang 应用就能在保持高并发的同时,与关系型数据库稳定协作。查询结果的内存开销与缓存策略可延伸阅读 https://plumephp.com/erlang-ets-caching/,慢查询与连接池指标的观测可参考 https://plumephp.com/erlang-logging-telemetry-observability/。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。