跳过正文
  1. 工作笔记/
  2. 数据集成/

OAuth2 大批量数据接口外部导入方案

·445 字·3 分钟
数据集成 OAuth2 Go MySQL 批量导入 接口同步
作者
molefool
目录

场景
#

对方提供的是 OAuth2 鉴权的数据接口,数据量很大。系统内置导入方式不适合直接处理全量数据,容易遇到超时、页面卡死、字段映射难排查、失败后不好续查的问题。

处理思路是单独写一个外部导入程序:

  • 先通过 OAuth2 client_credentials 获取 access_token
  • 再按接口返回的 cursor 分页拉取全量数据。
  • 每页解析成结构化模型,按固定字段映射到目标表。
  • 使用普通多值 INSERT 分批写入 MySQL。
  • 先支持 dry-run,只拉取和解析,不写库。
  • 失败时停止,不做覆盖、不做忽略、不自动清表。

程序结构
#

cmd/importer/main.go        # CLI 参数、HTTP client、数据库连接
internal/importer/api.go    # OAuth2 token、分页接口、重试和 token 刷新
internal/importer/model.go  # 接口 JSON 到入库字段的模型转换
internal/importer/runner.go # cursor 循环、日志、统计和停止条件
internal/importer/writer.go # 批量 INSERT 写入 MySQL

配置入口
#

敏感信息从命令行参数传入,不写死在代码里。

type Config struct {
    TokenURL        string
    DataURL         string
    ClientID        string
    ClientSecret    string
    BatchNo         string
    RequestTimeout  time.Duration
    Retries         int
    InsertBatchSize int
}

BatchNo 在一次全量导入中保持不变,方便接口侧和日志侧追踪同一批数据。

OAuth2 取 token
#

func (c *APIClient) RefreshToken(ctx context.Context) error {
    form := url.Values{}
    form.Set("grant_type", "client_credentials")
    form.Set("client_id", c.cfg.ClientID)
    form.Set("client_secret", c.cfg.ClientSecret)

    req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.cfg.TokenURL, strings.NewReader(form.Encode()))
    if err != nil {
        return err
    }
    req.Header.Set("Content-Type", "application/x-www-form-urlencoded")

    var token struct {
        AccessToken string `json:"access_token"`
    }
    if err := c.doJSON(req, &token); err != nil {
        return err
    }
    if token.AccessToken == "" {
        return fmt.Errorf("token response missing access_token")
    }
    c.token = token.AccessToken
    return nil
}

如果分页接口返回 401403,程序重新取 token 后再请求一次当前页。

分页拉取
#

func (c *APIClient) FetchPage(ctx context.Context, cursor string) (APIResponse, error) {
    if c.token == "" {
        if err := c.RefreshToken(ctx); err != nil {
            return APIResponse{}, err
        }
    }

    page, err := c.fetchPageWithToken(ctx, cursor)
    if err == nil {
        return page, nil
    }
    if !isAuthError(err) {
        return APIResponse{}, err
    }

    if err := c.RefreshToken(ctx); err != nil {
        return APIResponse{}, err
    }
    return c.fetchPageWithToken(ctx, cursor)
}

分页请求只带两个关键参数:

  • batchNo:本次导入批次号。
  • cursor:上一页返回的下一页游标,第一页为空。

如果接口返回 hasMore=true 但没有 nextCursor,直接报错停止,避免无限循环或漏数据。

主循环
#

func (r Runner) Run(ctx context.Context, opts RunOptions) error {
    var cursor string
    var totalInserted int64

    for pageNo := 1; ; pageNo++ {
        page, err := r.api.FetchPage(ctx, cursor)
        if err != nil {
            return fmt.Errorf("fetch page=%d cursor=%q inserted=%d: %w", pageNo, cursor, totalInserted, err)
        }

        inserted, err := r.writer.Insert(ctx, page.Data.Items)
        if err != nil {
            return fmt.Errorf("insert page=%d cursor=%q inserted=%d pageItems=%d: %w",
                pageNo, cursor, totalInserted, len(page.Data.Items), err)
        }
        totalInserted += inserted

        if !page.Data.HasMore {
            return nil
        }
        cursor = page.Data.NextCursor
    }
}

日志里记录页码、当前页数量、累计插入量、hasMorenextCursor、耗时和接口追踪号。出错时至少能定位到第几页、哪个游标、已插入多少。

批量写入
#

写库只做普通 INSERT,不做 UPDATEREPLACEINSERT IGNOREON DUPLICATE KEY UPDATE

func BuildInsertSQL(table string, rows []Row) (string, []any) {
    cols := InsertColumns()
    placeholders := "(" + strings.TrimRight(strings.Repeat("?,", len(cols)), ",") + ")"

    var b strings.Builder
    b.WriteString("INSERT INTO ")
    b.WriteString(quoteIdent(table))
    b.WriteString(" (")
    for i, col := range cols {
        if i > 0 {
            b.WriteString(", ")
        }
        b.WriteString(quoteIdent(col))
    }
    b.WriteString(") VALUES ")

    args := make([]any, 0, len(rows)*len(cols))
    for i, row := range rows {
        if i > 0 {
            b.WriteString(",")
        }
        b.WriteString(placeholders)
        args = append(args, row.Values()...)
    }
    return b.String(), args
}

表名和字段名只允许字母、数字和下划线,再统一加反引号,避免把不可信标识符拼进 SQL。

试跑方式
#

正式导入前先 dry-run,只拉一页并打印第一条数据的字段名,不打印字段值。

./data-import \
  --dry-run \
  --limit-pages 1 \
  --print-first-item-fields \
  --batch-no 202607020001 \
  --client-id 'CLIENT_ID' \
  --client-secret 'CLIENT_SECRET'

确认字段和解析逻辑无误后,再传数据库参数正式写入。批量大小可以先从 5001000 开始,根据接口和数据库压力调整。

注意点
#

  • 不要把接口地址、client secret、数据库密码写进仓库。
  • 目标表如果要清空,人工执行并确认,不让程序自动清表。
  • 遇到重复主键让数据库报错并停止,不静默忽略。
  • 对数字和日期字段要兼容空字符串、null、字符串数字和多种日期格式。
  • 大批量导入一定要保留批次号、页码、游标和累计数量日志。