You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

执行查询后更新时出现‘conn busy’错误的技术求助

解决pgx单连接下查询时执行更新出现「conn busy」错误

问题背景

需求为查询会员数据的同时,将对应记录的lastlogin字段更新为当前时间。使用Go语言结合pgx连接PostgreSQL时,出现「2023/08/18 15:07:21 conn busy」错误,查询功能正常但更新操作失败,调整更新代码位置、增加defer close语句后仍未解决。

原因分析

  1. 单连接资源冲突:当前使用单个pgx.Conn连接,在rows.Next()循环中,连接正被查询结果集rows占用,此时尝试用同一连接执行Exec更新操作,会触发连接忙的错误——pgx单个连接同一时间仅能处理一个数据库操作。
  2. 错误的延迟关闭位置:将defer rows.Close()放在for rows.Next()循环内部,会导致每次循环都注册一个延迟关闭操作,既冗余又可能提前关闭结果集,导致后续迭代出错。
  3. 时间处理不当:手动将时间格式化为字符串传入数据库,不如直接传递time.Time类型,pgx可自动处理与PostgreSQL时间类型的映射。

解决方案

方案1:改用pgx连接池(推荐)

单个连接无法同时处理多个操作,改用连接池可自动分配空闲连接处理更新请求,避免资源冲突。

修改数据库连接代码

package DBConnection

import (
    "context"
    "log"

    "github.com/jackc/pgx/v4/pgxpool"
)

var dbPool *pgxpool.Pool
var ctx = context.Background()

func InitDB() {
    poolConfig, err := pgxpool.ParseConfig("postgres://postgres:12345@localhost:5432/crd_audit?sslmode=disable")
    if err != nil {
        log.Fatalf("解析数据库连接配置失败: %v", err)
    }

    dbPool, err = pgxpool.ConnectConfig(ctx, poolConfig)
    if err != nil {
        log.Fatalf("无法连接到数据库池: %v", err)
    }
}

func GetDBPool() *pgxpool.Pool {
    return dbPool
}

修改业务代码

// SELECT MEMBER
func SearchMember(c *fiber.Ctx) error {
    var request Models.Request
    if err := c.BodyParser(&request); err != nil {
        return err
    }

    dbPool := DBConnection.GetDBPool()
    rows, err := dbPool.Query(ctx, "SELECT id, login, firstname, mi, lastname, email, passwd, institution, area, internet, level, lastlogin, region, loginToken, staffId, mobileNumber, userStat, isLogin FROM members WHERE login = $1 AND passwd = $2", request.Login, request.Password)
    if err != nil {
        log.Println(err)
        return c.Status(fiber.StatusInternalServerError).SendString("Internal Server Error")
    }
    // 将延迟关闭放在循环外,确保结果集迭代完成后再释放连接
    defer rows.Close()

    var Members []Models.SelectMembers
    var updateIDs []int

    for rows.Next() {
        var Member Models.SelectMembers
        var Institution sql.NullString
        var Area sql.NullString
        var Internet sql.NullString
        var Level sql.NullInt32
        var LastLogin sql.NullTime // 改用sql.NullTime处理时间类型
        var Region sql.NullInt32
        var LoginToken sql.NullString
        var StaffID sql.NullString
        var UserStatus sql.NullInt32
        var ISLogin sql.NullInt32

        err := rows.Scan(&Member.MembersID, &Member.LogIn, &Member.FullName.FirstName, &Member.FullName.MiddleName, &Member.FullName.LastName, &Member.Email, &Member.Password, &Institution, &Area, &Internet, &Level, &LastLogin, &Region, &LoginToken, &StaffID, &Member.MobileNumber, &UserStatus, &ISLogin)
        if err != nil {
            log.Println(err)
            continue
        }

        Member.Institution = Institution.String
        Member.Area = Area.String
        Member.Internet = Internet.String
        Member.Level = int(Level.Int32)
        if LastLogin.Valid {
            Member.LastLogin = LastLogin.Time.Format("2006-01-02 15:04:05")
            updateIDs = append(updateIDs, Member.MembersID)
        }
        Member.LoginToken = LoginToken.String
        Member.Region = int(Region.Int32)
        Member.StaffID = StaffID.String
        Member.UserStatus = int(UserStatus.Int32)
        Member.ISLogin = int(ISLogin.Int32)

        Member.Name = fmt.Sprintf("%s %s %s", Member.FullName.FirstName, Member.FullName.MiddleName, Member.FullName.LastName)
        Members = append(Members, Member)
    }

    // 批量更新lastlogin字段
    if len(updateIDs) > 0 {
        currentTime := time.Now()
        for _, id := range updateIDs {
            _, updateErr := dbPool.Exec(ctx, "UPDATE members SET lastlogin = $1 WHERE id = $2", currentTime, id)
            if updateErr != nil {
                log.Println(updateErr)
            }
        }
    }

    return c.Status(fiber.StatusOK).JSON(Members)
}

方案2:先读取所有查询结果,再执行更新(无需连接池)

若暂时不想改用连接池,可先将所有会员数据读取到内存中,关闭结果集后再执行更新操作,避免连接被占用:

// SELECT MEMBER
func SearchMember(c *fiber.Ctx) error {
    var request Models.Request
    if err := c.BodyParser(&request); err != nil {
        return err
    }

    db := DBConnection.GetDB()
    rows, err := db.Query(ctx, "SELECT id, login, firstname, mi, lastname, email, passwd, institution, area, internet, level, lastlogin, region, loginToken, staffId, mobileNumber, userStat, isLogin FROM members WHERE login = $1 AND passwd = $2", request.Login, request.Password)
    if err != nil {
        log.Println(err)
        return c.Status(fiber.StatusInternalServerError).SendString("Internal Server Error")
    }
    defer rows.Close()

    var Members []Models.SelectMembers
    var updateIDs []int

    // 先读取所有结果到内存
    for rows.Next() {
        var Member Models.SelectMembers
        var Institution sql.NullString
        var Area sql.NullString
        var Internet sql.NullString
        var Level sql.NullInt32
        var LastLogin sql.NullTime
        var Region sql.NullInt32
        var LoginToken sql.NullString
        var StaffID sql.NullString
        var UserStatus sql.NullInt32
        var ISLogin sql.NullInt32

        err := rows.Scan(&Member.MembersID, &Member.LogIn, &Member.FullName.FirstName, &Member.FullName.MiddleName, &Member.FullName.LastName, &Member.Email, &Member.Password, &Institution, &Area, &Internet, &Level, &LastLogin, &Region, &LoginToken, &StaffID, &Member.MobileNumber, &UserStatus, &ISLogin)
        if err != nil {
            log.Println(err)
            continue
        }

        Member.Institution = Institution.String
        Member.Area = Area.String
        Member.Internet = Internet.String
        Member.Level = int(Level.Int32)
        if LastLogin.Valid {
            Member.LastLogin = LastLogin.Time.Format("2006-01-02 15:04:05")
            updateIDs = append(updateIDs, Member.MembersID)
        }
        Member.LoginToken = LoginToken.String
        Member.Region = int(Region.Int32)
        Member.StaffID = StaffID.String
        Member.UserStatus = int(UserStatus.Int32)
        Member.ISLogin = int(ISLogin.Int32)

        Member.Name = fmt.Sprintf("%s %s %s", Member.FullName.FirstName, Member.FullName.MiddleName, Member.FullName.LastName)
        Members = append(Members, Member)
    }

    // 结果集已关闭,现在可用同一连接执行更新
    if len(updateIDs) > 0 {
        currentTime := time.Now()
        for _, id := range updateIDs {
            _, updateErr := db.Exec(ctx, "UPDATE members SET lastlogin = $1 WHERE id = $2", currentTime, id)
            if updateErr != nil {
                log.Println(updateErr)
            }
        }
    }

    return c.Status(fiber.StatusOK).JSON(Members)
}

额外优化点

  • 将lastlogin字段的处理改为sql.NullTime,避免字符串转换的麻烦和格式错误。
  • 更新操作直接传入time.Now(),pgx会自动适配PostgreSQL的timestamp类型。
  • 批量更新可进一步优化为使用IN子句减少数据库交互(需注意PostgreSQL中IN的参数数量限制)。

内容的提问来源于stack exchange,提问作者CJV

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.13 07:00:56