Erlang 数据库集成与连接池实战

系统讲解 Erlang 访问关系型数据库的完整方案:epgsql 与 mysql-otp 驱动选型、poolboy 连接池设计、事务与预处理语句、迁移与索引策略、故障重连与超时治理。

Erlang 应用访问关系型数据库的方式与其他语言有本质差异:没有 ORM 的对象映射负担,没有线程本地连接的概念,连接是进程持有的资源。这带来两个后果——连接池必须显式管理(否则进程数一涨就会打爆数据库),事务必须绑定在单个进程与单条连接上(否则跨进程共享连接会串事务)。理解这两点,是用好 Erlang 数据库集成的关键。本文将以 PostgreSQL 的 epgsql 与 MySQL 的 mysql-otp 为主线,讲解驱动选型、poolboy 连接池、事务与预处理语句、迁移与索引策略,以及超时与故障重连的治理。

一、Erlang 数据库访问生态

1.1 生态概览

层次方案说明
原生驱动epgsql、mysql-otp纯 Erlang,直接讲协议
查询构建pgo、sqerl提供池化与查询封装
Elixir ORMEctoElixir 侧的事实标准
连接池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 驱动能力对比

能力epgsqlmysql-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 + trigrampg_trgm 扩展
JSONB 字段GINjsonb_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_timeout5s
查询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/。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「erlang」更多文章

  1. rebar3 构建与发布:Erlang 工程化的完整工具箱
  2. RabbitMQ 与消息中间件:AMQP 模型与可靠投递
  3. Dialyzer 与类型规范:Erlang 静态分析的工程实践