fix: 稳定顺运宝当天分页同步 (#192)
This commit is contained in:
+143
-73
@@ -37,9 +37,10 @@ func mappingContextVersion(c repository.SybOrderContext) string {
|
||||
}
|
||||
|
||||
const (
|
||||
dateLayout = "2006-01-02"
|
||||
maxSpecifiedSyncDays = 31
|
||||
defaultSybMaxMatches = 10000
|
||||
dateLayout = "2006-01-02"
|
||||
maxSpecifiedSyncDays = 31
|
||||
defaultSybMaxMatches = 10000
|
||||
sybTodayListMaxAttempts = 3
|
||||
)
|
||||
|
||||
// SybSyncOptions 表示一次同步的日期选择。页面必须明确传入两端日期;两端都空
|
||||
@@ -149,44 +150,17 @@ type SybSyncDefaults struct {
|
||||
Warning string
|
||||
}
|
||||
|
||||
// DefaultSybSyncRange 通常预填最近三天;如果覆盖游标落后,则优先从游标
|
||||
// 当天连续补齐。缺口超过 31 天时只选最早一段,成功后下一次继续。
|
||||
func DefaultSybSyncRange(db *sql.DB, now time.Time) (SybSyncDefaults, error) {
|
||||
// DefaultSybSyncRange 固定预填昨天到今天。覆盖游标仍由同步服务维护,但不再
|
||||
// 改写采购员眼前的日期选择;需要补历史缺口时由采购员明确选择日期范围。
|
||||
func DefaultSybSyncRange(_ *sql.DB, now time.Time) (SybSyncDefaults, error) {
|
||||
today, err := time.Parse(dateLayout, dateOf(now))
|
||||
if err != nil {
|
||||
return SybSyncDefaults{}, fmt.Errorf("计算顺运宝今天日期失败: %w", err)
|
||||
}
|
||||
recentStart := today.AddDate(0, 0, -2)
|
||||
defaults := SybSyncDefaults{From: recentStart.Format(dateLayout), To: today.Format(dateLayout)}
|
||||
|
||||
lastSyncedAt, found, err := repository.GetSybLastSyncedAt(db)
|
||||
if err != nil {
|
||||
return SybSyncDefaults{}, err
|
||||
}
|
||||
if !found {
|
||||
return defaults, nil
|
||||
}
|
||||
lastTime, ok := model.ParseISO(lastSyncedAt)
|
||||
if !ok {
|
||||
return SybSyncDefaults{}, fmt.Errorf("上次同步时间 %q 解析失败", lastSyncedAt)
|
||||
}
|
||||
coveredDate, err := time.Parse(dateLayout, dateOf(lastTime))
|
||||
if err != nil || !coveredDate.Before(recentStart) {
|
||||
return defaults, nil
|
||||
}
|
||||
|
||||
segmentEnd := coveredDate.AddDate(0, 0, maxSpecifiedSyncDays-1)
|
||||
defaults.From = coveredDate.Format(dateLayout)
|
||||
if segmentEnd.Before(today) {
|
||||
defaults.To = segmentEnd.Format(dateLayout)
|
||||
defaults.Warning = fmt.Sprintf(
|
||||
"存在较长的未同步区间,已先选择最早 %d 天(%s 至 %s);本段成功后请继续同步下一段。",
|
||||
maxSpecifiedSyncDays, defaults.From, defaults.To)
|
||||
return defaults, nil
|
||||
}
|
||||
defaults.To = today.Format(dateLayout)
|
||||
defaults.Warning = fmt.Sprintf("检测到上次同步停在 %s,已自动扩展开始日期以补齐缺口。", defaults.From)
|
||||
return defaults, nil
|
||||
return SybSyncDefaults{
|
||||
From: today.AddDate(0, 0, -1).Format(dateLayout),
|
||||
To: today.Format(dateLayout),
|
||||
}, nil
|
||||
}
|
||||
|
||||
// cursorISOForDate 把“已连续覆盖到哪一天”保存成 UTC ISO。下一次仍从该日
|
||||
@@ -482,6 +456,77 @@ func roundYuanToCent(yuan float64) int64 {
|
||||
return int64(math.Round(yuan * 100))
|
||||
}
|
||||
|
||||
// sybDailyListResult 是一次单日列表分页尝试取得的数据。每次重试都重新创建,
|
||||
// 防止把前一次已发生 offset 位移的 ID 混进下一次。
|
||||
type sybDailyListResult struct {
|
||||
stockByID map[int64]syb.StockRow
|
||||
orderedIDs []int64
|
||||
observedTotal int
|
||||
}
|
||||
|
||||
// sybListDriftError 表示请求本身成功,但分页期间的列表没有形成稳定快照。
|
||||
// 只有今天允许有限重试;历史日期遇到它仍然立即失败。
|
||||
type sybListDriftError struct {
|
||||
message string
|
||||
}
|
||||
|
||||
func (e *sybListDriftError) Error() string { return e.message }
|
||||
|
||||
// loadSybDailyList 拉取一天的全部列表,并核对页长、分页前后总数及唯一 ID。
|
||||
// 发生快照漂移时仍返回本次已经取得的合法 ID,供今天最后一次尝试在完整拉取
|
||||
// 明细后安全 upsert;调用方不得因此把本次同步标记为成功或推进游标。
|
||||
func loadSybDailyList(ctx context.Context, client *syb.Client, date string, pageSize, expectedTotal int) (sybDailyListResult, error) {
|
||||
result := sybDailyListResult{
|
||||
stockByID: make(map[int64]syb.StockRow, expectedTotal),
|
||||
orderedIDs: make([]int64, 0, expectedTotal),
|
||||
observedTotal: expectedTotal,
|
||||
}
|
||||
for start := 0; start < expectedTotal; start += pageSize {
|
||||
pageIndex := start/pageSize + 1
|
||||
rows, pageCount, err := client.ListPage(ctx, date, date, start, pageIndex, pageSize)
|
||||
if err != nil {
|
||||
return result, fmt.Errorf("拉取 %s 货运单列表第 %d 页失败(已获取 %d/%d 张,本次同步整体作废,"+
|
||||
"下次会从同一个起始日期重新拉,靠 upsert 幂等不会重复计数): %w",
|
||||
date, pageIndex, len(result.orderedIDs), expectedTotal, err)
|
||||
}
|
||||
for _, row := range rows {
|
||||
if _, duplicate := result.stockByID[row.ID]; duplicate {
|
||||
continue
|
||||
}
|
||||
result.stockByID[row.ID] = row
|
||||
result.orderedIDs = append(result.orderedIDs, row.ID)
|
||||
}
|
||||
|
||||
expectedPageCount := min(pageSize, expectedTotal-start)
|
||||
if pageCount != expectedPageCount || len(rows) != expectedPageCount {
|
||||
afterTotal, totalErr := client.ListTotal(ctx, date, date, pageSize)
|
||||
if totalErr != nil {
|
||||
return result, fmt.Errorf("分页异常后重新查询 %s 货运单总数失败: %w", date, totalErr)
|
||||
}
|
||||
result.observedTotal = afterTotal
|
||||
return result, &sybListDriftError{message: fmt.Sprintf(
|
||||
"%s 货运单列表第 %d 页不完整:预期 %d 行,实际 %d 行",
|
||||
date, pageIndex, expectedPageCount, len(rows))}
|
||||
}
|
||||
}
|
||||
|
||||
afterTotal, err := client.ListTotal(ctx, date, date, pageSize)
|
||||
if err != nil {
|
||||
return result, fmt.Errorf("分页后重新查询 %s 货运单总数失败: %w", date, err)
|
||||
}
|
||||
result.observedTotal = afterTotal
|
||||
if afterTotal != expectedTotal {
|
||||
return result, &sybListDriftError{message: fmt.Sprintf(
|
||||
"%s 货运单总数在分页期间从 %d 变为 %d", date, expectedTotal, afterTotal)}
|
||||
}
|
||||
if len(result.orderedIDs) != expectedTotal {
|
||||
return result, &sybListDriftError{message: fmt.Sprintf(
|
||||
"%s 货运单列表不完整:预期 %d 张,分页后只有 %d 个唯一 ID",
|
||||
date, expectedTotal, len(result.orderedIDs))}
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
// RunSybSync 执行一次完整的顺运宝货运单同步。
|
||||
//
|
||||
// `[必须]` 调用方负责:
|
||||
@@ -575,7 +620,12 @@ func RunSybSyncWithOptions(ctx context.Context, db *sql.DB, client *syb.Client,
|
||||
}
|
||||
plans := make([]dailyPlan, 0, len(dates))
|
||||
allTotal := 0
|
||||
containsToday := false
|
||||
today := dateOf(now)
|
||||
for _, date := range dates {
|
||||
if date == today {
|
||||
containsToday = true
|
||||
}
|
||||
total, totalErr := client.ListTotal(ctx, date, date, pageSize)
|
||||
if totalErr != nil {
|
||||
report.Err = fmt.Errorf("查询 %s 货运单总数失败: %w", date, totalErr)
|
||||
@@ -597,7 +647,9 @@ func RunSybSyncWithOptions(ctx context.Context, db *sql.DB, client *syb.Client,
|
||||
allTotal += total
|
||||
plans = append(plans, dailyPlan{date: date, total: total})
|
||||
}
|
||||
if allTotal == 0 {
|
||||
// 历史范围全为 0 时可以直接结束。范围包含今天时仍再做一次列表快照
|
||||
// 校验,避免预检刚返回 0 就有新单进入而被当成完整成功。
|
||||
if allTotal == 0 && !containsToday {
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
if advanceCursor {
|
||||
if err := repository.SetSybLastSyncedAt(db, cursorAt); err != nil {
|
||||
@@ -620,54 +672,64 @@ func RunSybSyncWithOptions(ctx context.Context, db *sql.DB, client *syb.Client,
|
||||
// 的),但"这次同步整体算成功"这件事不能发生,否则漏掉的单永远补不回来。
|
||||
const detailBatch = 100
|
||||
for _, plan := range plans {
|
||||
if plan.total == 0 {
|
||||
if plan.total == 0 && plan.date != today {
|
||||
continue
|
||||
}
|
||||
stockByID := map[int64]syb.StockRow{}
|
||||
orderedIDs := make([]int64, 0, plan.total)
|
||||
for start := 0; start < plan.total; start += pageSize {
|
||||
pageIndex := start/pageSize + 1
|
||||
rows, pageCount, 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)
|
||||
|
||||
attempts := 1
|
||||
if plan.date == today {
|
||||
attempts = sybTodayListMaxAttempts
|
||||
}
|
||||
var listResult sybDailyListResult
|
||||
var unstableTodayErr error
|
||||
for attempt := 1; attempt <= attempts; attempt++ {
|
||||
listResult, err = loadSybDailyList(ctx, client, plan.date, pageSize, plan.total)
|
||||
if err == nil {
|
||||
unstableTodayErr = nil
|
||||
break
|
||||
}
|
||||
var driftErr *sybListDriftError
|
||||
if plan.date != today || !errors.As(err, &driftErr) {
|
||||
report.Err = fmt.Errorf("%w;本次同步停止且不推进游标", err)
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
expectedPageCount := min(pageSize, plan.total-start)
|
||||
if pageCount != expectedPageCount || len(rows) != expectedPageCount {
|
||||
report.Err = fmt.Errorf("%s 货运单列表第 %d 页不完整:预期 %d 行,实际 %d 行,"+
|
||||
"本次同步停止且不推进游标", plan.date, pageIndex, expectedPageCount, len(rows))
|
||||
|
||||
unstableTodayErr = driftErr
|
||||
if attempt == attempts {
|
||||
break
|
||||
}
|
||||
if listResult.observedTotal < 0 {
|
||||
report.Err = fmt.Errorf("重新查询 %s 货运单总数返回负数 %d",
|
||||
plan.date, listResult.observedTotal)
|
||||
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)
|
||||
otherDatesTotal := allTotal - plan.total
|
||||
if listResult.observedTotal > maxMatches-otherDatesTotal {
|
||||
report.Err = fmt.Errorf(
|
||||
"今天货运单变化后日期范围 %s ~ %s 的总数超过单次同步上限 %d,"+
|
||||
"已停止自动重试,请缩小日期范围或调大 max_matches",
|
||||
from, to, maxMatches)
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
allTotal = otherDatesTotal + listResult.observedTotal
|
||||
plan.total = listResult.observedTotal
|
||||
}
|
||||
afterTotal, totalErr := client.ListTotal(ctx, plan.date, plan.date, pageSize)
|
||||
if totalErr != nil {
|
||||
report.Err = fmt.Errorf("分页后重新查询 %s 货运单总数失败: %w", plan.date, totalErr)
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
if afterTotal != plan.total {
|
||||
report.Err = fmt.Errorf("%s 货运单总数在分页期间从 %d 变为 %d,"+
|
||||
"为防止 offset 分页漏单,本次同步停止且不推进游标", plan.date, plan.total, afterTotal)
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
if len(orderedIDs) != plan.total {
|
||||
report.Err = fmt.Errorf("%s 货运单列表不完整:预期 %d 张,分页后只有 %d 个唯一 ID,"+
|
||||
"本次同步停止且不推进游标", plan.date, plan.total, len(orderedIDs))
|
||||
otherDatesTotal := allTotal - plan.total
|
||||
largestTodayCount := max(listResult.observedTotal, len(listResult.orderedIDs))
|
||||
if largestTodayCount < 0 || largestTodayCount > maxMatches-otherDatesTotal {
|
||||
report.Err = fmt.Errorf(
|
||||
"今天货运单变化后日期范围 %s ~ %s 的总数超过单次同步上限 %d,"+
|
||||
"未保存本次超限数据,请缩小日期范围或调大 max_matches",
|
||||
from, to, maxMatches)
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
|
||||
stockByID := listResult.stockByID
|
||||
orderedIDs := listResult.orderedIDs
|
||||
report.StockCount += len(orderedIDs)
|
||||
|
||||
for i := 0; i < len(orderedIDs); i += detailBatch {
|
||||
@@ -698,6 +760,14 @@ func RunSybSyncWithOptions(ctx context.Context, db *sql.DB, client *syb.Client,
|
||||
}
|
||||
}
|
||||
}
|
||||
if unstableTodayErr != nil {
|
||||
report.Err = fmt.Errorf(
|
||||
"%s 当天货运单在连续 %d 次分页期间仍有变化;已保存最后一次取得的 %d 张货运单完整明细,"+
|
||||
"本次未形成稳定快照且不推进游标,下次同步将继续覆盖当天:%w",
|
||||
plan.date, sybTodayListMaxAttempts, len(orderedIDs), unstableTodayErr)
|
||||
report.FinishedAt = time.Now().UTC()
|
||||
return report
|
||||
}
|
||||
}
|
||||
|
||||
// ④ 全部成功,才更新 last_synced_at。
|
||||
|
||||
Reference in New Issue
Block a user