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
相关产品推荐
相关产品推荐

