NiFi中用Python脚本批量调用艺术品API并合并结果至单文件
NiFi ExecuteScript Python实现批量调用艺术品详情API
需求说明
输入文件包含存储所有艺术品ID的data数组,需基于这些ID批量调用https://api.article.edu/api/v1/artwork/{id}接口获取单条艺术品详情,最终将所有详情输出到单个JSON文件。
实现脚本
以下是可直接在NiFi ExecuteScript处理器中使用的Python代码:
import json import sys import requests from java.io import InputStreamReader, BufferedReader, OutputStreamWriter # 读取NiFi传入的输入流(包含ID数组的JSON) input_reader = BufferedReader(InputStreamReader(sys.stdin)) input_content = input_reader.readLine() input_json = json.loads(input_content) artwork_ids = input_json.get('data', []) # 存储所有成功获取的艺术品详情 artwork_details = [] # 批量调用详情API base_api_url = "https://api.article.edu/api/v1/artwork/" for art_id in artwork_ids: try: # 拼接单个艺术品详情的API地址 api_url = f"{base_api_url}{art_id}" # 发送GET请求 response = requests.get(api_url) # 检查HTTP响应状态,非200则抛出异常 response.raise_for_status() # 解析响应JSON并加入结果列表 artwork_details.append(response.json()) except Exception as err: # 将错误信息输出到NiFi日志(sys.stderr对应NiFi的处理器日志) print(f"获取艺术品ID {art_id} 详情失败: {str(err)}", file=sys.stderr) # 将结果写入输出流,NiFi会将此输出作为后续文件内容 output_writer = OutputStreamWriter(sys.stdout) output_writer.write(json.dumps(artwork_details, indent=2, ensure_ascii=False)) output_writer.flush() output_writer.close()
关键细节说明
- 输入/输出流处理:NiFi通过
sys.stdin传递输入文件内容,通过sys.stdout接收脚本输出内容,这里用Java IO类处理流的读取和写入,适配NiFi的运行环境。 - 异常处理:针对单个API请求的失败(如网络错误、HTTP 4xx/5xx)做捕获处理,避免单个请求失败导致整个脚本终止,错误信息会输出到NiFi的处理器日志中。
- JSON序列化:最终结果以格式化的JSON输出,方便后续NiFi处理器解析或直接生成文件。
替代方案(无requests库时)
如果NiFi环境未安装requests库,可使用Python标准库urllib.request替代,核心调用部分修改为:
import urllib.request import json for art_id in artwork_ids: try: api_url = f"{base_api_url}{art_id}" with urllib.request.urlopen(api_url) as response: # 读取并解码响应内容 response_content = response.read().decode('utf-8') artwork_details.append(json.loads(response_content)) except Exception as err: print(f"获取艺术品ID {art_id} 详情失败: {str(err)}", file=sys.stderr)
性能优化(大数量ID场景)
如果需要处理的艺术品ID数量较多,可使用线程池并发调用API提升效率(注意控制并发数,避免触发API速率限制):
from concurrent.futures import ThreadPoolExecutor def fetch_single_artwork(art_id): try: api_url = f"{base_api_url}{art_id}" response = requests.get(api_url) response.raise_for_status() return response.json() except Exception as err: print(f"获取艺术品ID {art_id} 详情失败: {str(err)}", file=sys.stderr) return None # 初始化线程池,设置最大并发数为5(可根据API限制调整) with ThreadPoolExecutor(max_workers=5) as executor: # 批量提交请求 results = executor.map(fetch_single_artwork, artwork_ids) # 过滤掉获取失败的结果 artwork_details = [res for res in results if res is not None]
内容的提问来源于stack exchange,提问作者vyshnavi dama
相关产品推荐
相关产品推荐

