Files
tyapi-server/internal/domains/api/services/processors/dwbg/dwbg6a2c_collect.go
2026-07-08 21:28:08 +08:00

143 lines
3.4 KiB
Go

package dwbg
import (
"context"
"encoding/json"
"fmt"
"sync"
"tyapi-server/internal/domains/api/dto"
"tyapi-server/internal/domains/api/services/processors"
"go.uber.org/zap"
)
type sinanAPICallInfo struct {
apiCode string
params map[string]interface{}
}
type sinanProcessorResult struct {
apiCode string
data interface{}
err error
}
// collectSinanAPIData 并发调用 9 个子接口;单接口失败仅记日志,对应 key 为 nil。
func collectSinanAPIData(ctx context.Context, params dto.DWBG6A2CReq, deps *processors.ProcessorDependencies, log *zap.Logger) map[string]interface{} {
apiCalls := []sinanAPICallInfo{
{apiCode: "YYSY35TA", params: map[string]interface{}{"mobile_no": params.MobileNo}},
{apiCode: "YYSYE7V5", params: map[string]interface{}{"mobile_no": params.MobileNo}},
{
apiCode: "YYSYH6F3",
params: map[string]interface{}{
"name": params.Name,
"id_card": params.IDCard,
"mobile_no": params.MobileNo,
},
},
{apiCode: "YYSYP0T4", params: map[string]interface{}{"mobile_no": params.MobileNo}},
{
apiCode: "FLXGDEA9",
params: map[string]interface{}{
"name": params.Name,
"id_card": params.IDCard,
"authorized": "1",
},
},
{
apiCode: "FLXG8B4D",
params: map[string]interface{}{
"mobile_no": params.MobileNo,
"authorized": "1",
},
},
{
apiCode: "FLXG7E8F",
params: map[string]interface{}{
"name": params.Name,
"id_card": params.IDCard,
"mobile_no": params.MobileNo,
},
},
{
apiCode: "JRZQ5E9F",
params: map[string]interface{}{
"name": params.Name,
"id_card": params.IDCard,
"mobile_no": params.MobileNo,
"authorized": "1",
},
},
{
apiCode: "JRZQ1D09",
params: map[string]interface{}{
"name": params.Name,
"id_card": params.IDCard,
"mobile_no": params.MobileNo,
"authorized": "1",
},
},
}
apiData := make(map[string]interface{}, len(apiCalls))
results := make(chan sinanProcessorResult, len(apiCalls))
var wg sync.WaitGroup
for _, apiCall := range apiCalls {
wg.Add(1)
go func(ac sinanAPICallInfo) {
defer wg.Done()
defer func() {
if r := recover(); r != nil {
log.Error("调用司南子接口时发生panic",
zap.String("api_code", ac.apiCode),
zap.Any("panic", r),
)
results <- sinanProcessorResult{apiCode: ac.apiCode, err: fmt.Errorf("处理器panic: %v", r)}
}
}()
paramsBytes, err := json.Marshal(ac.params)
if err != nil {
log.Warn("序列化司南子接口参数失败",
zap.String("api_code", ac.apiCode),
zap.Error(err),
)
results <- sinanProcessorResult{apiCode: ac.apiCode, err: err}
return
}
data, err := callProcessor(ctx, ac.apiCode, paramsBytes, deps)
results <- sinanProcessorResult{apiCode: ac.apiCode, data: data, err: err}
}(apiCall)
}
go func() {
wg.Wait()
close(results)
}()
successCount := 0
for result := range results {
if result.err != nil {
log.Warn("调用司南子接口失败,将使用默认值",
zap.String("api_code", result.apiCode),
zap.Error(result.err),
)
apiData[result.apiCode] = nil
continue
}
apiData[result.apiCode] = result.data
successCount++
}
log.Info("司南子接口调用完成",
zap.Int("total", len(apiCalls)),
zap.Int("success", successCount),
zap.Int("failed", len(apiCalls)-successCount),
)
return apiData
}