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

PySpark如何通过翻译API返回的JSON结果构造多语言新DataFrame

PySpark 对接翻译API生成多语种数据集实现方案

前置配置

先定义全局通用参数,避免硬编码:

from pyspark.sql import functions as F
from pyspark.sql.types import ArrayType, StringType
import requests
from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type

# 自定义配置项
TARGET_LANGS = {"zh", "en", "ja", "ko"} # 按需替换为你的目标语种集合
TRANSLATE_API_URL = "你的翻译API接口地址"
API_KEY = "你的接口密钥"
MAX_RETRY_TIMES = 3 # 接口调用最大重试次数

核心实现步骤

步骤1:计算待翻译语种列表

通过UDF判断每行需要翻译的目标语种,过滤不需要翻译的行:

@F.udf(returnType=ArrayType(StringType()))
def get_pending_langs(current_lang: str) -> list:
    return list(TARGET_LANGS - {current_lang})

# 原始DataFrame打标,区分原始数据和翻译数据
raw_df = df.withColumn("is_translated", F.lit(False))\
           .withColumn("pending_langs", get_pending_langs(F.col("col3")))

步骤2:封装API调用逻辑

加重试机制处理接口超时、限流等异常场景:

@retry(
    stop=stop_after_attempt(MAX_RETRY_TIMES),
    wait=wait_exponential(multiplier=1, min=2, max=10),
    retry=retry_if_exception_type((requests.exceptions.RequestException, ValueError))
)
def request_translate(text: str, target_langs: list) -> dict:
    req_body = {
        "text": text,
        "target_langs": target_langs
    }
    headers = {"Authorization": f"Bearer {API_KEY}"}
    resp = requests.post(
        TRANSLATE_API_URL,
        json=req_body,
        headers=headers,
        timeout=15
    )
    resp.raise_for_status()
    # 按需调整解析逻辑,和API返回结构对齐
    return resp.json()

步骤3:分区处理生成结果集

用mapPartitions按分区处理数据,减少客户端初始化开销,同时方便控制请求速率:

def process_partition(rows):
    for row in rows:
        row_data = row.asDict()
        pending_langs = row_data.pop("pending_langs")
        # 无待翻译语种直接返回原始行
        if not pending_langs:
            yield row_data
            continue
        # 调用翻译接口
        translate_res = request_translate(row_data["col2"], pending_langs)
        # 返回原始行
        yield row_data
        # 生成翻译后的新行
        for lang, translated_text in translate_res.items():
            new_row = row_data.copy()
            new_row["col2"] = translated_text
            new_row["col3"] = lang
            new_row["is_translated"] = True
            yield new_row

# RDD处理后转DataFrame
result_rdd = raw_df.rdd.mapPartitions(process_partition)
result_df = spark.createDataFrame(result_rdd, schema=raw_df.drop("pending_langs").schema)

生产环境优化建议

  • 翻译结果去重:提前对相同col2内容做去重,相同文本只调用一次API,翻译完成后再关联回原始数据,可大幅降低接口调用量
  • 限流控制:如果API有QPS限制,可在分区处理逻辑中添加请求间隔控制,避免触发限流规则,日调用数千次的量级只需简单的sleep控制即可满足要求
  • 异常兜底:添加降级逻辑,超过重试次数的请求可将失败数据写入异常表,后续定时重试,避免任务整体失败
  • 缓存复用:可将历史翻译结果存入Redis或Hive表,调用API前优先查询缓存,进一步减少无效请求
  • 批量请求优化:如果翻译API支持多文本批量提交,可在分区内攒N条文本后统一调用接口,提升处理效率

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 06:09:02