如何使用Pyspark/Python从Microsoft Graph API高效获取组织层级结构
实现方案
1 先调整表结构
先给default.orgTree新增两个必要字段,用来关联父子节点和标记层级:
parent_id:存储当前用户直接上级的用户IDlevel:存储当前节点的层级,CEO下属为1层,向下依次递增
2 初始化第一层数据
改造你现有的初始调用代码,写入基础数据时补充新增字段:
import requests import pyspark.sql.functions as f from pyspark.sql.types import * # 替换为实际的CEO用户ID、Graph API访问令牌 ceo_id = "替换为你的CEO用户ID" headers = {"Authorization": "Bearer 你的Graph API访问令牌"} # 拉取CEO的直接下属 url = f"https://graph.microsoft.com/v1.0/users/{ceo_id}/directReports" response = requests.get(url, headers=headers) df = spark.read.json(sc.parallelize([response.text])) df = df.select(f.explode_outer('value').alias('data')) df = df.select("data.*") # 补充父ID、层级字段 df = df.withColumn("parent_id", f.lit(ceo_id))\ .withColumn("level", f.lit(1)) # 首次覆盖写入表 df.write.mode("overwrite").saveAsTable("default.orgTree")
3 迭代拉取全量层级数据
用广度优先遍历逻辑循环拉取下层数据,直到没有新的下属返回就终止迭代:
current_level = 1 while True: # 取出当前层所有需要查询下属的用户ID current_level_ids = [ row.id for row in spark.sql( f"select distinct id from default.orgTree where level = {current_level}" ).collect() ] if not current_level_ids: break # 没有当前层用户,结束循环 next_level_data = [] for user_id in current_level_ids: url = f"https://graph.microsoft.com/v1.0/users/{user_id}/directReports" try: response = requests.get(url, headers=headers, timeout=10) response.raise_for_status() res_json = response.json() direct_reports = res_json.get("value", []) if direct_reports: # 给每个下属补充父ID、层级字段 for report in direct_reports: report["parent_id"] = user_id report["level"] = current_level + 1 next_level_data.append(report) except Exception as e: # 可自行打印异常信息排查问题,比如权限不足、限流、用户不存在等 print(f"用户{user_id}下属查询失败: {str(e)}") continue if not next_level_data: break # 没有下一层数据,结束循环 # 对齐字段后追加写入表 next_level_df = spark.createDataFrame(next_level_data) existing_cols = spark.table("default.orgTree").columns # 补全缺失字段,过滤多余字段 for col in existing_cols: if col not in next_level_df.columns: next_level_df = next_level_df.withColumn(col, f.lit(None)) next_level_df = next_level_df.select(existing_cols) # 追加写入 next_level_df.write.mode("append").saveAsTable("default.orgTree") current_level += 1
4 效率优化建议
- 若组织用户量较大,可将同步for循环调用API改为异步请求(用Python的aiohttp库实现),能大幅降低接口调用耗时
- 要处理Graph API的限流规则,若返回429状态码,需按响应头
Retry-After的值等待后重试 - 提前做好访问令牌的自动刷新逻辑,避免迭代过程中令牌过期导致任务失败
- 可额外增加数据去重逻辑,避免存在多汇报线的用户被重复写入表
内容的提问来源于stack exchange,提问作者user16714516
相关产品推荐
相关产品推荐

