场景#
对方提供的是 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
}如果分页接口返回 401 或 403,程序重新取 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
}
}日志里记录页码、当前页数量、累计插入量、hasMore、nextCursor、耗时和接口追踪号。出错时至少能定位到第几页、哪个游标、已插入多少。
批量写入#
写库只做普通 INSERT,不做 UPDATE、REPLACE、INSERT IGNORE 或 ON 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'确认字段和解析逻辑无误后,再传数据库参数正式写入。批量大小可以先从 500 或 1000 开始,根据接口和数据库压力调整。
注意点#
- 不要把接口地址、client secret、数据库密码写进仓库。
- 目标表如果要清空,人工执行并确认,不让程序自动清表。
- 遇到重复主键让数据库报错并停止,不静默忽略。
- 对数字和日期字段要兼容空字符串、
null、字符串数字和多种日期格式。 - 大批量导入一定要保留批次号、页码、游标和累计数量日志。
