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 }