Get Bee messages API增量存储JSON至数据湖:日期排序与每日刷新优化
解决Get Bee Messages API增量存储的日期序列与每日刷新问题
问题诊断
你的代码存在以下核心问题导致日期序列异常和重复数据:
- 固定
start参数:每次运行都使用同一时间戳,重复拉取历史数据,最终出现重复日期条目 - 循环缩进错误:API调用、数据处理逻辑不在
for bid in bid_list循环内,仅会处理最后一个BID - BID列表格式错误:
bid_list是包含逗号分隔字符串的单元素列表,而非单个BID的数组,API无法正确识别多个设备ID - 无数据去重与排序:直接追加新数据,未去除重复条目,也未按日期排序保证序列正确性
调整方案
要实现每日刷新并维护最新日期→最早日期的序列,需要:
- 动态计算
start参数,每日拉取最新1天的数据(或按需拉取最近10天但去重) - 修正循环与BID列表格式,确保每个BID都被正确处理
- 合并数据时去重,再按日期字段倒序排序后保存
- 提前创建目标目录,避免文件操作报错
完整修正代码
import requests import json import os from datetime import datetime, timedelta url = "https://api.roambee.com/bee/bee_messages" dir_path = "/dbfs/jhgfd/adata/dfs/roambee/raw/api_responses_get_bee_messages/" # 修正BID列表:拆分为单个设备ID元素 bid_list = ["3456", "12345", "2345", "345"] headers = { "apikey": "rtyrewrERTYUIOLIKJHGFRDSASDERFTGYHUJKILJHGFtyujhtrewsadfgh", "Content-Type": "application/json" } # 提前创建目录,避免后续文件操作报错 if not os.path.exists(dir_path): os.makedirs(dir_path) # 动态计算start时间:每日运行时拉取前一天的完整数据 end_date = datetime.now() start_date = end_date - timedelta(days=1) # 转换为API要求的毫秒级时间戳 start_timestamp = int(start_date.timestamp() * 1000) for bid in bid_list: params = { "bid": bid, "start": str(start_timestamp), "days": 1 # 每日拉取1天数据,避免重复存储历史数据 } # 调用API并捕获请求错误 response = requests.get(url, headers=headers, params=params) response.raise_for_status() new_bee_data = response.json() response_file_path = os.path.join(dir_path, f"response_data_{bid}.json") combined_data = [] # 加载现有数据(如果文件存在) if os.path.exists(response_file_path): with open(response_file_path, "r") as response_file: combined_data = json.load(response_file) # 数据去重:用消息时间戳作为唯一标识,避免重复条目 existing_timestamps = {item['messageTimestamp'] for item in combined_data} for item in new_bee_data: if item['messageTimestamp'] not in existing_timestamps: combined_data.append(item) # 按时间戳倒序排序,保证最新数据排在最前面(对应理想序列11,10,9...1) # 可替换为你转换后的人类可读日期字段,比如'formattedDate' combined_data.sort(key=lambda x: x['messageTimestamp'], reverse=True) # 保存合并后的数据 with open(response_file_path, "w") as response_file: json.dump(combined_data, response_file, indent=2) print(f"Response data for BID {bid} updated in:", response_file_path) # 测试读取Delta表(可选) delta_table_path = "/mnt/analyticsdata/rms/roambee/raw/api_responses_get_bee_messages/response_data_868199058105959.json" first_bee_test = spark.read.format("json").load(delta_table_path) display(first_bee_test)
关键说明
- 动态时间计算:使用
datetime模块生成每日起始时间戳,确保每次运行只拉取最新1天的数据,从根源避免重复 - 数据去重:通过集合存储现有数据的时间戳,过滤重复的消息条目
- 排序逻辑:按时间戳倒序排序,保证最新数据排在最前面,符合你需要的日期序列
- 错误处理:添加
response.raise_for_status()捕获API请求失败的情况,便于排查问题 - 灵活调整拉取范围:如果需要每次拉取最近10天的数据,只需修改时间计算部分:
# 拉取最近10天的数据 start_date = datetime.now() - timedelta(days=10) start_timestamp = int(start_date.timestamp() * 1000) params["days"] = 10
该场景下仍需保留去重逻辑,避免重复存储历史数据。
内容的提问来源于stack exchange,提问作者vish
相关产品推荐
相关产品推荐

