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

基于Spark并行化优化API调用:批量客户详情查询改造

基于Spark并行调用客户API的实现方案

需求背景

现有两个API接口:

  • get_all_customers():获取所有客户ID列表
  • get_specifc_customer_details():获取指定客户详情,当前为串行调用,需改为Spark并行调用,最终生成包含customer_id、cust_status_1、cust_status_2三列的DataFrame。

原代码问题

原get_specifc_customer_details()通过循环串行调用API,无法利用Spark的分布式计算能力,客户量较大时效率极低。

改造实现方案

核心思路

  1. 先通过get_all_customers()获取客户ID列表,转为Spark DataFrame(每行对应一个客户ID)
  2. 定义用户自定义函数(UDF),针对单个客户ID调用API并返回状态数据
  3. 在DataFrame上调用UDF,并行获取每个客户的详情,解析后生成目标列

完整代码实现

import requests
from pyspark.sql import SparkSession
from pyspark.sql.functions import udf, col
from pyspark.sql.types import StructType, StructField, StringType, IntegerType  # 按需调整数据类型

# 初始化SparkSession(未初始化时执行)
spark = SparkSession.builder.appName("CustomerAPIParallelCall").getOrCreate()

# 复用请求头,避免重复创建
API_HEADERS = {
    'Content-Type': 'application/json; charset=utf-8',
    'Authorization': 'Bearer Token'
}

def get_all_customers():
    url = "https://test.url.com/api/v1/customers?status=Customer"
    response = requests.request("GET", url, headers=API_HEADERS)
    # 直接提取客户ID列表
    return [cust["uid"] for cust in response.json()["customers"]]

# 单个客户API调用函数,带异常处理
def fetch_customer_status(customer_id):
    try:
        url = f"https://test.url.com/api/v1/customers/{customer_id}"
        response = requests.request("GET", url, headers=API_HEADERS)
        response.raise_for_status()  # 捕获HTTP错误
        cust_data = response.json()["customer"]
        return (cust_data["status_1"], cust_data["status_2"])
    except Exception as e:
        # 异常时返回默认值,避免任务失败
        print(f"获取客户{customer_id}详情失败: {str(e)}")
        return (None, None)

# 注册UDF,指定返回类型(按需调整)
status_schema = StructType([
    StructField("cust_status_1", StringType(), nullable=True),
    StructField("cust_status_2", StringType(), nullable=True)
])
fetch_status_udf = udf(fetch_customer_status, status_schema)

# 主流程
if __name__ == "__main__":
    # 1. 获取客户ID并转为DataFrame
    customer_ids = get_all_customers()
    customer_df = spark.createDataFrame(customer_ids, StringType()).toDF("customer_id")
    
    # 2. 并行调用API,解析结果生成目标列
    final_df = customer_df.withColumn("status_data", fetch_status_udf(col("customer_id"))) \
        .select(
            col("customer_id"),
            col("status_data.cust_status_1").alias("cust_status_1"),
            col("status_data.cust_status_2").alias("cust_status_2")
        )
    
    # 查看结果
    final_df.show()
    # 按需保存结果
    # final_df.write.mode("overwrite").parquet("/path/to/save")

关键说明

  • 并行执行:Spark会将DataFrame分区分配到不同节点,各节点任务并行调用API,大幅提升效率
  • 请求头复用:将headers定义为全局常量,避免每次调用API重复创建
  • 异常处理:UDF中加入异常捕获,单个API调用失败不会导致整个任务终止,同时打印错误信息便于排查
  • 数据类型适配:根据API返回的status_1和status_2实际类型调整UDF的返回Schema(比如数字类型用IntegerType)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 11:16:16