fix: 增强顺运宝同步完整性保护 (#58)
This commit is contained in:
+148
-50
@@ -30,7 +30,11 @@ import (
|
||||
"cmautobuy/admin/syb"
|
||||
)
|
||||
|
||||
const dateLayout = "2006-01-02"
|
||||
const (
|
||||
dateLayout = "2006-01-02"
|
||||
maxSpecifiedSyncDays = 31
|
||||
defaultSybMaxMatches = 10000
|
||||
)
|
||||
|
||||
// SybSyncOptions 表示一次同步的日期选择。From/To 都为空是日常自动增量;
|
||||
// 两者都有值是操作员明确发起的指定日期补同步。
|
||||
@@ -72,9 +76,35 @@ func NewSybSyncOptions(fromRaw, toRaw string, now time.Time) (SybSyncOptions, er
|
||||
if to > dateOf(now) {
|
||||
return SybSyncOptions{}, fmt.Errorf("结束日期不能晚于顺运宝服务端今天(%s)", dateOf(now))
|
||||
}
|
||||
days := int(toDate.Sub(fromDate).Hours()/24) + 1
|
||||
if days > maxSpecifiedSyncDays {
|
||||
return SybSyncOptions{}, fmt.Errorf("指定日期同步最多选择 %d 天,当前范围为 %d 天,请分段同步",
|
||||
maxSpecifiedSyncDays, days)
|
||||
}
|
||||
return SybSyncOptions{From: from, To: to}, nil
|
||||
}
|
||||
|
||||
// splitDateRange 把闭区间拆成逐日查询段。顺运宝只支持日期粒度,按天拉取
|
||||
// 能把 offset 分页期间新增数据造成的位移限制在当天。
|
||||
func splitDateRange(from, to string) ([]string, error) {
|
||||
fromDate, err := time.Parse(dateLayout, from)
|
||||
if err != nil || fromDate.Format(dateLayout) != from {
|
||||
return nil, fmt.Errorf("开始日期格式不正确,请使用 YYYY-MM-DD")
|
||||
}
|
||||
toDate, err := time.Parse(dateLayout, to)
|
||||
if err != nil || toDate.Format(dateLayout) != to {
|
||||
return nil, fmt.Errorf("结束日期格式不正确,请使用 YYYY-MM-DD")
|
||||
}
|
||||
if fromDate.After(toDate) {
|
||||
return nil, fmt.Errorf("开始日期不能晚于结束日期")
|
||||
}
|
||||
dates := make([]string, 0, int(toDate.Sub(fromDate).Hours()/24)+1)
|
||||
for day := fromDate; !day.After(toDate); day = day.AddDate(0, 0, 1) {
|
||||
dates = append(dates, day.Format(dateLayout))
|
||||
}
|
||||
return dates, nil
|
||||
}
|
||||
|
||||
// sybLocation 是顺运宝服务端的时区。
|
||||
//
|
||||
// `[必须]` 日期范围筛的是服务端的 created,而它是 UTC+8——实测
|
||||
@@ -292,7 +322,7 @@ func RunSybSyncWithOptions(ctx context.Context, db *sql.DB, client *syb.Client,
|
||||
}
|
||||
maxMatches := cfg.MaxMatches
|
||||
if maxMatches <= 0 {
|
||||
maxMatches = 500
|
||||
maxMatches = defaultSybMaxMatches
|
||||
}
|
||||
|
||||
lastSyncedAt, _, err := repository.GetSybLastSyncedAt(db)
|
||||
@@ -327,20 +357,41 @@ func RunSybSyncWithOptions(ctx context.Context, db *sql.DB, client *syb.Client,
|
||||
}
|
||||
report.From, report.To = from, to
|
||||
|
||||
total, err := client.ListTotal(ctx, from, to, pageSize)
|
||||
dates, err := splitDateRange(from, to)
|
||||
if err != nil {
|
||||
report.Err = fmt.Errorf("查询货运单总数失败: %w", err)
|
||||
report.Err = err
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
if total > maxMatches {
|
||||
report.Err = fmt.Errorf(
|
||||
"日期范围 %s ~ %s 内有 %d 张货运单,超过单次同步上限 %d,"+
|
||||
"请缩小日期范围或联系维护者调大 max_matches", from, to, total, maxMatches)
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
type dailyPlan struct {
|
||||
date string
|
||||
total int
|
||||
}
|
||||
if total == 0 {
|
||||
plans := make([]dailyPlan, 0, len(dates))
|
||||
allTotal := 0
|
||||
for _, date := range dates {
|
||||
total, totalErr := client.ListTotal(ctx, date, date, pageSize)
|
||||
if totalErr != nil {
|
||||
report.Err = fmt.Errorf("查询 %s 货运单总数失败: %w", date, totalErr)
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
if total < 0 {
|
||||
report.Err = fmt.Errorf("查询 %s 货运单总数返回负数 %d", date, total)
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
if total > maxMatches-allTotal {
|
||||
report.Err = fmt.Errorf(
|
||||
"日期范围 %s ~ %s 内货运单总数超过单次同步上限 %d,"+
|
||||
"请缩小日期范围或在确认机器和网络容量后调大 max_matches", from, to, maxMatches)
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
allTotal += total
|
||||
plans = append(plans, dailyPlan{date: date, total: total})
|
||||
}
|
||||
if allTotal == 0 {
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
if advanceCursor {
|
||||
if err := repository.SetSybLastSyncedAt(db, model.NowISO()); err != nil {
|
||||
@@ -350,30 +401,7 @@ func RunSybSyncWithOptions(ctx context.Context, db *sql.DB, client *syb.Client,
|
||||
return report
|
||||
}
|
||||
|
||||
// ① 翻页拉全部货运单行。
|
||||
stockByID := map[int64]syb.StockRow{}
|
||||
var orderedIDs []int64
|
||||
for start := 0; start < total; start += pageSize {
|
||||
pageIndex := start/pageSize + 1
|
||||
rows, err := client.ListPage(ctx, from, to, start, pageIndex, pageSize)
|
||||
if err != nil {
|
||||
report.Err = fmt.Errorf("拉取货运单列表第 %d 页失败(已获取 %d/%d 张,本次同步整体作废,"+
|
||||
"下次会从同一个起始日期重新拉,靠 upsert 幂等不会重复计数): %w",
|
||||
pageIndex, len(orderedIDs), total, err)
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
for _, row := range rows {
|
||||
if _, dup := stockByID[row.ID]; dup {
|
||||
continue
|
||||
}
|
||||
stockByID[row.ID] = row
|
||||
orderedIDs = append(orderedIDs, row.ID)
|
||||
}
|
||||
}
|
||||
report.StockCount = len(orderedIDs)
|
||||
|
||||
// ② 按 100 个一批取明细,③ 逐张货运单写库。
|
||||
// ① 每天单独翻页,② 按 100 个一批取明细,③ 逐张货运单写库。
|
||||
//
|
||||
// `[必须]` 不放在一个大事务里——几千条明细的事务会长时间持锁;
|
||||
// 按货运单为单位提交,失败了已成功的部分保留(下次重拉会 upsert
|
||||
@@ -383,29 +411,71 @@ func RunSybSyncWithOptions(ctx context.Context, db *sql.DB, client *syb.Client,
|
||||
// 不更新 last_synced_at——已经成功写入的部分不回滚(它们本身是幂等
|
||||
// 的),但"这次同步整体算成功"这件事不能发生,否则漏掉的单永远补不回来。
|
||||
const detailBatch = 100
|
||||
for i := 0; i < len(orderedIDs); i += detailBatch {
|
||||
end := i + detailBatch
|
||||
if end > len(orderedIDs) {
|
||||
end = len(orderedIDs)
|
||||
for _, plan := range plans {
|
||||
if plan.total == 0 {
|
||||
continue
|
||||
}
|
||||
batch := orderedIDs[i:end]
|
||||
|
||||
details, err := client.DetailListByStock(ctx, batch)
|
||||
if err != nil {
|
||||
report.Err = fmt.Errorf("拉取货运单明细失败(本次同步整体作废,"+
|
||||
"已写入的数据保留,下次重拉会 upsert 覆盖): %w", err)
|
||||
stockByID := map[int64]syb.StockRow{}
|
||||
orderedIDs := make([]int64, 0, plan.total)
|
||||
for start := 0; start < plan.total; start += pageSize {
|
||||
pageIndex := start/pageSize + 1
|
||||
rows, responseTotal, listErr := client.ListPage(ctx, plan.date, plan.date, start, pageIndex, pageSize)
|
||||
if listErr != nil {
|
||||
report.Err = fmt.Errorf("拉取 %s 货运单列表第 %d 页失败(已获取 %d/%d 张,本次同步整体作废,"+
|
||||
"下次会从同一个起始日期重新拉,靠 upsert 幂等不会重复计数): %w",
|
||||
plan.date, pageIndex, len(orderedIDs), plan.total, listErr)
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
if responseTotal != plan.total {
|
||||
report.Err = fmt.Errorf("%s 货运单总数在分页期间从 %d 变为 %d,"+
|
||||
"为防止 offset 分页漏单,本次同步停止且不推进游标", plan.date, plan.total, responseTotal)
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
for _, row := range rows {
|
||||
if _, dup := stockByID[row.ID]; dup {
|
||||
continue
|
||||
}
|
||||
stockByID[row.ID] = row
|
||||
orderedIDs = append(orderedIDs, row.ID)
|
||||
}
|
||||
}
|
||||
if len(orderedIDs) != plan.total {
|
||||
report.Err = fmt.Errorf("%s 货运单列表不完整:预期 %d 张,分页后只有 %d 个唯一 ID,"+
|
||||
"本次同步停止且不推进游标", plan.date, plan.total, len(orderedIDs))
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
report.StockCount += len(orderedIDs)
|
||||
|
||||
for _, d := range details {
|
||||
stockRow := stockByID[d.ID]
|
||||
if err := writeStockDetail(db, cfg.BaseURL, stockRow, d, &report); err != nil {
|
||||
report.Err = fmt.Errorf("写入货运单 %s(id=%d)失败(本次同步整体作废,"+
|
||||
"已写入的数据保留): %w", d.Code, d.ID, err)
|
||||
for i := 0; i < len(orderedIDs); i += detailBatch {
|
||||
end := i + detailBatch
|
||||
if end > len(orderedIDs) {
|
||||
end = len(orderedIDs)
|
||||
}
|
||||
batch := orderedIDs[i:end]
|
||||
details, detailErr := client.DetailListByStock(ctx, batch)
|
||||
if detailErr != nil {
|
||||
report.Err = fmt.Errorf("拉取 %s 货运单明细失败(本次同步整体作废,"+
|
||||
"已写入的数据保留,下次重拉会 upsert 覆盖): %w", plan.date, detailErr)
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
if completeErr := validateDetailBatch(batch, details); completeErr != nil {
|
||||
report.Err = fmt.Errorf("%s 货运单明细不完整:%w;本次同步停止且不推进游标", plan.date, completeErr)
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
for _, d := range details {
|
||||
stockRow := stockByID[d.ID]
|
||||
if err := writeStockDetail(db, cfg.BaseURL, stockRow, d, &report); err != nil {
|
||||
report.Err = fmt.Errorf("写入货运单 %s(id=%d)失败(本次同步整体作废,"+
|
||||
"已写入的数据保留): %w", d.Code, d.ID, err)
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -420,6 +490,34 @@ func RunSybSyncWithOptions(ctx context.Context, db *sql.DB, client *syb.Client,
|
||||
return report
|
||||
}
|
||||
|
||||
// validateDetailBatch 确认批量明细响应与请求 ID 一一对应。任何缺失、重复、
|
||||
// 意外 ID 或空商品明细都会让同步失败,避免在数据不完整时推进游标。
|
||||
func validateDetailBatch(requested []int64, details []syb.StockDetail) error {
|
||||
wanted := make(map[int64]struct{}, len(requested))
|
||||
for _, id := range requested {
|
||||
wanted[id] = struct{}{}
|
||||
}
|
||||
seen := make(map[int64]struct{}, len(details))
|
||||
for _, detail := range details {
|
||||
if _, ok := wanted[detail.ID]; !ok {
|
||||
return fmt.Errorf("响应包含未请求的货运单 id=%d", detail.ID)
|
||||
}
|
||||
if _, duplicate := seen[detail.ID]; duplicate {
|
||||
return fmt.Errorf("响应重复返回货运单 id=%d", detail.ID)
|
||||
}
|
||||
if len(detail.Details) == 0 {
|
||||
return fmt.Errorf("货运单 id=%d 没有返回商品明细", detail.ID)
|
||||
}
|
||||
seen[detail.ID] = struct{}{}
|
||||
}
|
||||
for _, id := range requested {
|
||||
if _, ok := seen[id]; !ok {
|
||||
return fmt.Errorf("响应缺少货运单 id=%d", id)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// writeStockDetail 把一张货运单的全部商品明细写进 syb_orders,
|
||||
// 一张货运单一个事务(工单 #46「按货运单为单位提交」)。
|
||||
func writeStockDetail(db *sql.DB, baseURL string, stockRow syb.StockRow, detail syb.StockDetail, report *SyncReport) error {
|
||||
|
||||
Reference in New Issue
Block a user