feat: 支持批量导入 PDD 商品链接 (#104)

This commit is contained in:
chengma
2026-08-10 16:14:24 +08:00
parent 7bd75510b2
commit 188725ff64
15 changed files with 911 additions and 32 deletions
+228
View File
@@ -0,0 +1,228 @@
// PDD 商品链接 Excel 批量导入。
//
// 本文件只负责文件校验、逐行解析、业务去重和事务编排;SQL 仍然全部在
// repository,HTTP 上传细节仍然全部在 handler/web。
package service
import (
"bytes"
"database/sql"
"errors"
"fmt"
"path/filepath"
"strings"
"github.com/xuri/excelize/v2"
"cmautobuy/admin/repository"
)
const (
// MaxPddUploadBytes 足够容纳 5000 条单列链接,同时限制误传大文件。
MaxPddUploadBytes = 10 * 1024 * 1024
// MaxPddImportRows 限制一次导入的非空数据行,避免误发超大批次。
MaxPddImportRows = 5000
pddLinkHeader = "拼多多链接"
)
var pddXLSXMagic = []byte{0x50, 0x4B, 0x03, 0x04}
// ErrInvalidPddImport 区分“文件需要采购员修正”和“服务器写库失败”。
// Handler 据此选择 400 或 500,避免把数据库内部错误直接显示到页面。
var ErrInvalidPddImport = errors.New("PDD 导入文件无效")
// IsInvalidPddImport 判断错误是否来自上传文件内容或结构。
func IsInvalidPddImport(err error) bool {
return errors.Is(err, ErrInvalidPddImport)
}
// PddImportFailure 是一条不能导入的 Excel 数据行。
type PddImportFailure struct {
Row int
Raw string
Reason string
}
// PddImportResult 是本次导入的完整统计。
// GoodsIDs 保存全部合法且去重后的商品,供页面显式创建本批采集任务。
type PddImportResult struct {
TotalRows int
CreatedCount int
ExistingCount int
RevivedCount int
DuplicateCount int
Failures []PddImportFailure
GoodsIDs []string
}
// ValidatePddUpload 在落盘和解析前校验上传文件的基本安全边界。
func ValidatePddUpload(filename string, size int64, head []byte) error {
if size <= 0 {
return fmt.Errorf("文件是空的")
}
if size > MaxPddUploadBytes {
return fmt.Errorf("文件超过 10MB 上限(当前约 %.1fMB)", float64(size)/1024/1024)
}
if !strings.EqualFold(filepath.Ext(filename), ".xlsx") {
return fmt.Errorf("只允许 .xlsx 文件,收到的是 %q", filepath.Ext(filename))
}
if !bytes.HasPrefix(head, pddXLSXMagic) {
return fmt.Errorf("文件内容不像一个 xlsx(zip 格式)文件,可能是改了扩展名的其他文件")
}
return nil
}
type pddImportEntry struct {
GoodsID string
URL string
}
// ImportPddExcel 导入第一个非空工作表的“拼多多链接”列。
//
// 格式错误只影响对应行,合法行仍会进入同一个数据库事务;数据库写入任一步
// 失败则整体回滚,避免留下无法解释的半批数据。
func ImportPddExcel(db *sql.DB, path string) (*PddImportResult, error) {
book, err := excelize.OpenFile(path)
if err != nil {
return nil, fmt.Errorf("%w:打开 Excel 失败(%v)", ErrInvalidPddImport, err)
}
defer book.Close()
result, entries, err := parsePddImportWorkbook(book)
if err != nil {
return nil, fmt.Errorf("%w:%v", ErrInvalidPddImport, err)
}
if len(entries) == 0 {
return result, nil
}
tx, err := db.Begin()
if err != nil {
return nil, fmt.Errorf("开始 PDD 导入事务失败: %w", err)
}
defer tx.Rollback()
for _, entry := range entries {
_, outcome, err := repository.EnsurePddProductWithOutcome(tx, entry.GoodsID, entry.URL)
if err != nil {
return nil, fmt.Errorf("写入 PDD 商品 %s 失败: %w", entry.GoodsID, err)
}
switch outcome {
case repository.PddProductCreated:
result.CreatedCount++
case repository.PddProductExisting:
result.ExistingCount++
case repository.PddProductRevived:
result.RevivedCount++
default:
return nil, fmt.Errorf("PDD 商品 %s 返回了未知导入结果 %q", entry.GoodsID, outcome)
}
}
if err := tx.Commit(); err != nil {
return nil, fmt.Errorf("提交 PDD 导入事务失败: %w", err)
}
return result, nil
}
func parsePddImportWorkbook(book *excelize.File) (*PddImportResult, []pddImportEntry, error) {
for _, sheet := range book.GetSheetList() {
rows, err := book.Rows(sheet)
if err != nil {
return nil, nil, fmt.Errorf("读取工作表 %q 失败: %w", sheet, err)
}
result, entries, found, parseErr := parsePddImportRows(rows, sheet)
closeErr := rows.Close()
if parseErr != nil {
return nil, nil, parseErr
}
if closeErr != nil {
return nil, nil, fmt.Errorf("关闭工作表 %q 失败: %w", sheet, closeErr)
}
if found {
return result, entries, nil
}
}
return nil, nil, fmt.Errorf("Excel 中没有非空工作表,请把链接填在第一列并使用表头“%s”", pddLinkHeader)
}
func parsePddImportRows(rows *excelize.Rows, sheet string) (*PddImportResult, []pddImportEntry, bool, error) {
result := &PddImportResult{}
entries := make([]pddImportEntry, 0, 64)
seen := make(map[string]struct{}, 64)
foundHeader := false
rowNumber := 0
for rows.Next() {
rowNumber++
columns, err := rows.Columns()
if err != nil {
return nil, nil, false, fmt.Errorf("读取工作表 %q 第 %d 行失败: %w", sheet, rowNumber, err)
}
if rowIsBlank(columns) {
continue
}
if !foundHeader {
foundHeader = true
header := firstColumn(columns)
if header != pddLinkHeader {
return nil, nil, false, fmt.Errorf(
"工作表 %q 第一个非空行的第一列应为“%s”,实际是 %q",
sheet, pddLinkHeader, header)
}
continue
}
result.TotalRows++
if result.TotalRows > MaxPddImportRows {
return nil, nil, false, fmt.Errorf(
"非空数据超过 %d 条上限(在工作表 %q 第 %d 行发现超限),请拆成多个文件导入",
MaxPddImportRows, sheet, rowNumber)
}
raw := firstColumn(columns)
if raw == "" {
result.Failures = append(result.Failures, PddImportFailure{
Row: rowNumber, Raw: "(空)", Reason: "第一列“拼多多链接”不能为空",
})
continue
}
goodsID, err := ParsePddGoodsID(raw)
if err != nil {
result.Failures = append(result.Failures, PddImportFailure{
Row: rowNumber, Raw: raw, Reason: err.Error(),
})
continue
}
if _, duplicate := seen[goodsID]; duplicate {
result.DuplicateCount++
continue
}
seen[goodsID] = struct{}{}
entries = append(entries, pddImportEntry{GoodsID: goodsID, URL: raw})
result.GoodsIDs = append(result.GoodsIDs, goodsID)
}
if err := rows.Error(); err != nil {
return nil, nil, false, fmt.Errorf("遍历工作表 %q 失败: %w", sheet, err)
}
if foundHeader && result.TotalRows == 0 {
return nil, nil, false, fmt.Errorf(
"工作表 %q 只有表头,下面没有 PDD 商品链接", sheet)
}
return result, entries, foundHeader, nil
}
func rowIsBlank(columns []string) bool {
for _, value := range columns {
if strings.TrimSpace(value) != "" {
return false
}
}
return true
}
func firstColumn(columns []string) string {
if len(columns) == 0 {
return ""
}
return strings.TrimSpace(columns[0])
}
+264
View File
@@ -0,0 +1,264 @@
package service
import (
"database/sql"
"fmt"
"path/filepath"
"strings"
"testing"
"github.com/xuri/excelize/v2"
"cmautobuy/admin/model"
"cmautobuy/admin/repository"
)
func TestValidatePddUpload(t *testing.T) {
validHead := []byte{0x50, 0x4B, 0x03, 0x04, 0x14}
for _, tc := range []struct {
name string
filename string
size int64
head []byte
wantErr bool
}{
{"合法文件", "pdd.xlsx", 1024, validHead, false},
{"扩展名大小写", "pdd.XLSX", 1024, validHead, false},
{"空文件", "pdd.xlsx", 0, validHead, true},
{"超过上限", "pdd.xlsx", MaxPddUploadBytes + 1, validHead, true},
{"错误扩展名", "pdd.xls", 1024, validHead, true},
{"伪造扩展名", "pdd.xlsx", 1024, []byte("not zip"), true},
} {
t.Run(tc.name, func(t *testing.T) {
err := ValidatePddUpload(tc.filename, tc.size, tc.head)
if (err != nil) != tc.wantErr {
t.Fatalf("ValidatePddUpload() err=%v, wantErr=%t", err, tc.wantErr)
}
})
}
}
func TestParsePddImportWorkbook_忽略空表并逐行反馈(t *testing.T) {
book := excelize.NewFile()
t.Cleanup(func() { book.Close() })
if err := book.SetSheetName("Sheet1", "空表"); err != nil {
t.Fatal(err)
}
dataSheet, err := book.NewSheet("链接")
if err != nil {
t.Fatal(err)
}
book.SetActiveSheet(dataSheet)
rows := [][]string{
{pddLinkHeader},
{"https://mobile.yangkeduo.com/goods.html?goods_id=737116531267"},
{"https://mobile.pinduoduo.com/goods.html?goods_id=647453710994"},
{"https://mobile.yangkeduo.com/goods.html?goods_id=737116531267&from=duplicate"},
{"https://p.pinduoduo.com/short"},
{"", "第二列误填内容"},
}
setWorkbookRows(t, book, "链接", rows)
result, entries, err := parsePddImportWorkbook(book)
if err != nil {
t.Fatalf("解析失败: %v", err)
}
if result.TotalRows != 5 || result.DuplicateCount != 1 || len(result.Failures) != 2 {
t.Fatalf("统计不对: %+v", result)
}
if len(entries) != 2 || len(result.GoodsIDs) != 2 {
t.Fatalf("合法商品应为 2 条: entries=%d ids=%v", len(entries), result.GoodsIDs)
}
if result.Failures[0].Row != 5 || result.Failures[1].Row != 6 {
t.Fatalf("失败行号不对: %+v", result.Failures)
}
if !strings.Contains(result.Failures[0].Reason, "goods_id") ||
!strings.Contains(result.Failures[1].Reason, "不能为空") {
t.Fatalf("失败原因不可操作: %+v", result.Failures)
}
}
func TestParsePddImportWorkbook_脱敏固定样本(t *testing.T) {
book, err := excelize.OpenFile(filepath.Join("..", "testdata", "pdd_links.xlsx"))
if err != nil {
t.Fatal(err)
}
defer book.Close()
result, entries, err := parsePddImportWorkbook(book)
if err != nil {
t.Fatalf("解析脱敏样本失败: %v", err)
}
if result.TotalRows != 5 || len(entries) != 3 || result.DuplicateCount != 1 || len(result.Failures) != 1 {
t.Fatalf("脱敏样本统计不对: result=%+v entries=%d", result, len(entries))
}
}
func TestParsePddImportWorkbook_错误表头和行数上限整体拒绝(t *testing.T) {
t.Run("错误表头", func(t *testing.T) {
book := excelize.NewFile()
defer book.Close()
setWorkbookRows(t, book, "Sheet1", [][]string{{"链接"}, {pddURL("737116531267")}})
_, _, err := parsePddImportWorkbook(book)
if err == nil || !strings.Contains(err.Error(), pddLinkHeader) {
t.Fatalf("应明确提示正确表头,实际 %v", err)
}
})
t.Run("只有表头", func(t *testing.T) {
book := excelize.NewFile()
defer book.Close()
setWorkbookRows(t, book, "Sheet1", [][]string{{pddLinkHeader}})
_, _, err := parsePddImportWorkbook(book)
if err == nil || !strings.Contains(err.Error(), "只有表头") {
t.Fatalf("应提示没有商品链接,实际 %v", err)
}
})
t.Run("超过5000条", func(t *testing.T) {
book := excelize.NewFile()
defer book.Close()
rows := make([][]string, 0, MaxPddImportRows+2)
rows = append(rows, []string{pddLinkHeader})
for i := 0; i <= MaxPddImportRows; i++ {
rows = append(rows, []string{pddURL(fmt.Sprintf("%012d", 100000000000+i))})
}
setWorkbookRows(t, book, "Sheet1", rows)
_, _, err := parsePddImportWorkbook(book)
if err == nil || !strings.Contains(err.Error(), "5000") || !strings.Contains(err.Error(), "拆成多个文件") {
t.Fatalf("超限提示不明确: %v", err)
}
})
}
func TestImportPddExcel_新建复用复活重复和失败可同时统计(t *testing.T) {
db := newPddImportSQLiteDB(t, false)
existingID := "647453710994"
revivedID := "753136429979"
if _, err := repository.EnsurePddProduct(db, existingID, pddURL(existingID)); err != nil {
t.Fatal(err)
}
if err := repository.SetCollectResult(db, existingID, "已采集商品", "测试店铺", `{"skus":[1]}`); err != nil {
t.Fatal(err)
}
if _, err := repository.EnsurePddProduct(db, revivedID, pddURL(revivedID)); err != nil {
t.Fatal(err)
}
if err := repository.SetCollectResult(db, revivedID, "删除前商品", "旧店铺", `{"skus":[1]}`); err != nil {
t.Fatal(err)
}
if _, err := repository.SoftDeletePddProducts(db, []string{revivedID}); err != nil {
t.Fatal(err)
}
path := filepath.Join("..", "testdata", "pdd_links.xlsx")
result, err := ImportPddExcel(db, path)
if err != nil {
t.Fatalf("导入失败: %v", err)
}
if result.TotalRows != 5 || result.CreatedCount != 1 || result.ExistingCount != 1 ||
result.RevivedCount != 1 || result.DuplicateCount != 1 || len(result.Failures) != 1 {
t.Fatalf("导入统计不对: %+v", result)
}
if len(result.GoodsIDs) != 3 {
t.Fatalf("本次导入集合应跨分页保留全部 3 个 ID,实际 %v", result.GoodsIDs)
}
existing, _ := repository.GetPddProductByGoodsID(db, existingID)
if existing.CollectStatus != model.CollectCollected || existing.SkusJSON == "" || existing.Title != "已采集商品" {
t.Fatalf("已存在商品的采集结果被覆盖: %+v", existing)
}
if !strings.Contains(existing.URL, "from=excel") {
t.Errorf("已存在商品应更新同 goods_id 的链接原文,实际 %q", existing.URL)
}
revived, _ := repository.GetPddProductByGoodsID(db, revivedID)
if revived.IsDeleted() || revived.CollectStatus != model.CollectPending || revived.SkusJSON != "" || revived.Title != "" {
t.Fatalf("复活商品未沿用现有清理规则: %+v", revived)
}
second, err := ImportPddExcel(db, path)
if err != nil {
t.Fatalf("重复导入失败: %v", err)
}
if second.CreatedCount != 0 || second.ExistingCount != 3 || second.RevivedCount != 0 {
t.Fatalf("重复导入应全部复用,实际 %+v", second)
}
var count int
if err := db.QueryRow(`SELECT COUNT(*) FROM pdd_products`).Scan(&count); err != nil || count != 3 {
t.Fatalf("重复导入产生重复商品: count=%d err=%v", count, err)
}
var taskCount int
if err := db.QueryRow(`SELECT COUNT(*) FROM tasks`).Scan(&taskCount); err != nil || taskCount != 0 {
t.Fatalf("上传不应自动创建采集任务: count=%d err=%v", taskCount, err)
}
}
func TestImportPddExcel_数据库中途失败时合法行整体回滚(t *testing.T) {
db := newPddImportSQLiteDB(t, true)
_, err := ImportPddExcel(db, filepath.Join("..", "testdata", "pdd_links.xlsx"))
if err == nil || IsInvalidPddImport(err) {
t.Fatalf("应返回数据库写入错误,实际 %v", err)
}
var count int
if err := db.QueryRow(`SELECT COUNT(*) FROM pdd_products`).Scan(&count); err != nil {
t.Fatal(err)
}
if count != 0 {
t.Fatalf("第一条合法商品没有随第二条失败回滚,仍有 %d 条", count)
}
}
func newPddImportSQLiteDB(t *testing.T, failSecondInsert bool) *sql.DB {
t.Helper()
dsn := fmt.Sprintf("file:pdd_import_%s?mode=memory&cache=shared", strings.ReplaceAll(t.Name(), "/", "_"))
db, err := sql.Open("sqlite", dsn)
if err != nil {
t.Fatal(err)
}
db.SetMaxOpenConns(1)
t.Cleanup(func() { db.Close() })
if _, err := db.Exec(`
CREATE TABLE pdd_products (
id INTEGER PRIMARY KEY AUTOINCREMENT,
goods_id TEXT NOT NULL UNIQUE,
url TEXT NOT NULL,
title TEXT, shop_name TEXT, skus_json TEXT,
collect_status TEXT NOT NULL,
collect_msg TEXT, artifact_ref TEXT, collected_at TEXT,
deleted_at TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL
);
CREATE TABLE tasks (task_id TEXT PRIMARY KEY);
`); err != nil {
t.Fatal(err)
}
if failSecondInsert {
if _, err := db.Exec(`
CREATE TRIGGER fail_second_pdd BEFORE INSERT ON pdd_products
WHEN NEW.goods_id = '647453710994'
BEGIN SELECT RAISE(ABORT, 'simulated write failure'); END;
`); err != nil {
t.Fatal(err)
}
}
return db
}
func setWorkbookRows(t *testing.T, book *excelize.File, sheet string, rows [][]string) {
t.Helper()
for rowIndex, row := range rows {
for columnIndex, value := range row {
cell, err := excelize.CoordinatesToCellName(columnIndex+1, rowIndex+1)
if err != nil {
t.Fatal(err)
}
if err := book.SetCellValue(sheet, cell, value); err != nil {
t.Fatal(err)
}
}
}
}
func pddURL(goodsID string) string {
return "https://mobile.yangkeduo.com/goods.html?goods_id=" + goodsID
}