You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Java 17下批量产品属性API调用的并发优化方案咨询

批量产品属性提交的并发优化方案(Java 17)

问题背景

我手里有一批关联了多个属性的产品,需要根据每个属性的attributeType调用外部API单独保存attributeValue(该外部API不支持批量提交属性)。目前采用同步遍历的实现:遍历批量产品列表→遍历单个产品→收集属性列表→遍历每个属性对象→判断attributeType→调用对应API(全程用forEach循环)。但每个产品至少包含50个属性,该实现处理耗时极长,现寻求Java 17环境下的并发优化方案。

相关实体类代码

BulkEditRequest类

import java.time.LocalDateTime;
import java.util.Map;
import java.util.Objects;

import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.Getter;
import lombok.Setter;
import lombok.ToString;

@Getter @Setter
@AllArgsConstructor
@ToString
@Data
public class BulkEditRequest{

    private String productId;
    private String productNum;
    private Map<String, AttributesResponse> attributeList;
}

AttributesResponse类

import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.Getter;
import lombok.Setter;
import lombok.ToString;

@Getter @Setter
@AllArgsConstructor
@ToString
@Data
public class AttributesResponse {
    private String attributeName;
    private String attributeValue;
    private String attributeType;
    private String attributeId;
    
    public AttributesResponse() {
        super();
    }
}

UploadEntity类

import lombok.AllArgsConstructor;
import lombok.Getter;
import lombok.NoArgsConstructor;
import lombok.Setter;

@Getter
@Setter
@AllArgsConstructor
@NoArgsConstructor
public class UploadEntity {

    // Price相关字段
    private String price;
    private String countryCode;
    private String uomID;

    // Identifier相关字段
    private String productIdentifier;

    // Attribute相关字段
    private String attributeId;
}

现有同步实现代码

import java.util.ArrayList;
import java.util.List;

import org.springframework.http.HttpEntity;
import org.springframework.http.ResponseEntity;
import org.springframework.stereotype.Service;
import org.springframework.web.client.RestTemplate;

@Service
public class ServiceImpl {

    private RestTemplate template;

    public ServiceImpl(RestTemplate template) {
        this.template = template;
    }

    public void saveProductsInFlex(List<BulkEditRequest> attributesByProd) {

        for (BulkEditRequest singleProduct : attributesByProd) {
            List<AttributesResponse> attributeList = new ArrayList<>(
                    singleProduct.getAttributeList().values());

            for (AttributesResponse attribute : attributeList) {

                if ("Price".equalsIgnoreCase(attribute.getAttributeType())) {
                    setProductPrice(singleProduct, attribute);
                } else if ("Identifier".equalsIgnoreCase(attribute.getAttributeType())) {
                    setProductIdentifier(singleProduct, attribute);
                } else if ("Attribute".equalsIgnoreCase(attribute.getAttributeType())) {
                    setProductAttribute(singleProduct, attribute);
                }
            }
        }
    }

    private void setProductPrice(BulkEditRequest singleProduct, AttributesResponse attributeResponse) {
        HttpEntity<UploadEntity> flexEntity = null;
        String uri = null;
        if (singleProduct != null) {
            uri = Constants.BASEURL + singleProduct.getProductNum() + Constants.PRICES;
        }
        try {
            UploadEntity entity = new UploadEntity();
            entity.setPrice(attributeResponse.getAttributeValue());
            entity.setUomID("1");
            entity.setCountryCode("US");
            flexEntity = new HttpEntity<>(entity);
            ResponseEntity<String> response = template.postForEntity(uri, flexEntity, String.class);
        } catch (Exception e) {
            // 原实现未处理异常
        }
    }

    private void setProductAttribute(BulkEditRequest singleProduct, AttributesResponse attributeResponse) {
        HttpEntity<UploadEntity> flexEntity = null;
        String uri = null;
        if (singleProduct != null) {
            uri = Constants.BASEURL + singleProduct.getProductNum() + Constants.ATTRIBUTES;
        }
        try {
            UploadEntity entity = new UploadEntity();
            entity.setAttributeId(attributeResponse.getAttributeId());
            flexEntity = new HttpEntity<>(entity);
            ResponseEntity<String> response = template.postForEntity(uri, flexEntity, String.class);
        } catch (Exception e) {
            // 原实现未处理异常
        }
    }

    private void setProductIdentifier(BulkEditRequest singleProduct, AttributesResponse attributeResponse) {
        HttpEntity<UploadEntity> flexEntity = null;
        String uri = null;
        if (singleProduct != null) {
            uri = Constants.BASEURL + singleProduct.getProductNum() + Constants.IDENTIFIERS;
        }
        try {
            UploadEntity entity = new UploadEntity();
            entity.setProductIdentifier(attributeResponse.getAttributeValue());
            flexEntity = new HttpEntity<>(entity);
            ResponseEntity<String> response = template.postForEntity(uri, flexEntity, String.class);
        } catch (Exception e) {
            // 原实现未处理异常
        }
    }
}

Java 17并发优化方案

核心优化思路

  1. 替换同步HTTP客户端为异步非阻塞实现:用WebClient替代RestTemplate,避免线程阻塞等待API响应,提升线程利用率。
  2. 并行处理属性提交:利用并行流或自定义线程池,并行处理每个产品的属性提交,充分利用多核CPU资源。
  3. 添加重试与异常处理:针对网络波动或API临时不可用场景添加重试机制,同时完善错误日志便于后续补偿。
  4. 控制并发数:自定义线程池限制并发请求数,避免压垮外部API。

优化后的代码实现

import org.springframework.stereotype.Service;
import org.springframework.web.reactive.function.client.WebClient;
import reactor.core.publisher.Mono;
import reactor.util.retry.Retry;

import java.time.Duration;
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

@Service
public class OptimizedServiceImpl {

    private final WebClient webClient;
    private final ExecutorService executorService;

    // 初始化WebClient和自定义线程池(并发数根据外部API限制调整)
    public OptimizedServiceImpl(WebClient.Builder webClientBuilder) {
        this.webClient = webClientBuilder.baseUrl(Constants.BASEURL).build();
        this.executorService = Executors.newFixedThreadPool(20);
    }

    public void saveProductsInFlex(List<BulkEditRequest> attributesByProd) {
        // 并行处理每个产品
        attributesByProd.parallelStream().forEach(singleProduct -> {
            List<AttributesResponse> attributeList = List.copyOf(singleProduct.getAttributeList().values());
            // 并行处理每个属性的API提交
            attributeList.parallelStream().forEach(attribute -> {
                switch (attribute.getAttributeType().toLowerCase()) {
                    case "price" -> submitPrice(singleProduct, attribute).subscribe();
                    case "identifier" -> submitIdentifier(singleProduct, attribute).subscribe();
                    case "attribute" -> submitAttribute(singleProduct, attribute).subscribe();
                }
            });
        });

        // 等待所有异步任务完成(如需同步等待结果)
        executorService.shutdown();
        try {
            if (!executorService.awaitTermination(10, java.util.concurrent.TimeUnit.MINUTES)) {
                executorService.shutdownNow();
            }
        } catch (InterruptedException e) {
            executorService.shutdownNow();
            Thread.currentThread().interrupt();
        }
    }

    // 异步提交Price属性
    private Mono<Void> submitPrice(BulkEditRequest singleProduct, AttributesResponse attributeResponse) {
        String uri = "/" + singleProduct.getProductNum() + Constants.PRICES;
        UploadEntity entity = new UploadEntity();
        entity.setPrice(attributeResponse.getAttributeValue());
        entity.setUomID("1");
        entity.setCountryCode("US");

        return webClient.post()
                .uri(uri)
                .bodyValue(entity)
                .retrieve()
                .bodyToMono(Void.class)
                // 指数退避重试,最多3次,间隔1秒起步
                .retryWhen(Retry.backoff(3, Duration.ofSeconds(1))
                        .filter(throwable -> throwable instanceof RuntimeException))
                .doOnError(e -> {
                    // 记录详细错误日志,便于后续排查
                    System.err.printf("提交Price失败 | 产品编号:%s | 属性值:%s | 错误信息:%s%n",
                            singleProduct.getProductNum(), attributeResponse.getAttributeValue(), e.getMessage());
                });
    }

    // 异步提交Identifier属性
    private Mono<Void> submitIdentifier(BulkEditRequest singleProduct, AttributesResponse attributeResponse) {
        String uri = "/" + singleProduct.getProductNum() + Constants.IDENTIFIERS;
        UploadEntity entity = new UploadEntity();
        entity.setProductIdentifier(attributeResponse.getAttributeValue());

        return webClient.post()
                .uri(uri)
                .bodyValue(entity)
                .retrieve()
                .bodyToMono(Void.class)
                .retryWhen(Retry.backoff(3, Duration.ofSeconds(1)))
                .doOnError(e -> {
                    System.err.printf("提交Identifier失败 | 产品编号:%s | 属性值:%s | 错误信息:%s%n",
                            singleProduct.getProductNum(), attributeResponse.getAttributeValue(), e.getMessage());
                });
    }

    // 异步提交Attribute属性
    private Mono<Void> submitAttribute(BulkEditRequest singleProduct, AttributesResponse attributeResponse) {
        String uri = "/" + singleProduct.getProductNum() + Constants.ATTRIBUTES;
        UploadEntity entity = new UploadEntity();
        entity.setAttributeId(attributeResponse.getAttributeId());

        return webClient.post()
                .uri(uri)
                .bodyValue(entity)
                .retrieve()
                .bodyToMono(Void.class)
                .retryWhen(Retry.backoff(3, Duration.ofSeconds(1)))
                .doOnError(e -> {
                    System.err.printf("提交Attribute失败 | 产品编号:%s | 属性ID:%s | 错误信息:%s%n",
                            singleProduct.getProductNum(), attributeResponse.getAttributeId(), e.getMessage());
                });
    }
}

额外优化建议

  • 配置WebClient连接池:通过HttpClient配置连接池大小,进一步提升HTTP请求的效率。
  • 批量失败补偿:将失败的请求记录到数据库或消息队列,后续通过定时任务重试。
  • 监控与告警:添加指标监控(如请求成功率、响应时间),针对异常情况设置告警。

内容的提问来源于stack exchange,提问作者Suprita Venkat

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 10:44:50