背景
对于特大文件上传到飞书的功能实现
直接贴源代码,文件 local_chunk_uploader.py

import xbot
from xbot import print, sleep
from . import package
from .package import variables as glv
import os
import requests
import zlib
import time
import json
from requests_toolbelt.multipart.encoder import MultipartEncoder
import shutil
class FeishuLocalChunkUploader:
API_ENDPOINTS = {
"material": {
"prepare": "https://open.feishu.cn/open-apis/drive/v1/material/upload_prepare",
"upload_part": "https://open.feishu.cn/open-apis/drive/v1/material/upload_part",
"finish": "https://open.feishu.cn/open-apis/drive/v1/material/upload_finish"
},
"file": {
"prepare": "https://open.feishu.cn/open-apis/drive/v1/files/upload_prepare",
"upload_part": "https://open.feishu.cn/open-apis/drive/v1/files/upload_part",
"finish": "https://open.feishu.cn/open-apis/drive/v1/files/upload_finish"
}
}
def __init__(self, app_id, app_secret, debug=False):
self.app_id = app_id
self.app_secret = app_secret
self.tenant_access_token = None
self.token_expire_time = 0
self.debug = debug
self.chunk_dir = None # 分片文件存储目录
def _log(self, message):
if self.debug:
print(f"[DEBUG] {time.strftime('%H:%M:%S')} - {message}")
def _get_tenant_token(self):
if self.tenant_access_token and time.time() < self.token_expire_time - 300:
return self.tenant_access_token
url = "https://open.feishu.cn/open-apis/auth/v3/tenant_access_token/internal"
headers = {"Content-Type": "application/json; charset=utf-8"}
payload = {"app_id": self.app_id, "app_secret": self.app_secret}
try:
response = requests.post(url, headers=headers, json=payload, timeout=10)
response.raise_for_status()
result = response.json()
if result.get("code") != 0:
raise Exception(f"token错误: {result.get('msg')}")
self.tenant_access_token = result["tenant_access_token"]
self.token_expire_time = time.time() + result.get("expire", 7200)
return self.tenant_access_token
except Exception as e:
raise Exception(f"获取token失败: {str(e)}")
def _prepare_upload(self, upload_type, file_name, parent_node, file_size, extra=None):
url = self.API_ENDPOINTS[upload_type]["prepare"]
token = self._get_tenant_token()
headers = {
"Authorization": f"Bearer {token}",
"Content-Type": "application/json; charset=utf-8"
}
payload = {
"file_name": file_name[:250],
"parent_type": "docx_image" if upload_type == "material" else "explorer",
"parent_node": parent_node,
"size": file_size
}
if extra:
payload["extra"] = extra
self._log(f"预上传参数: {payload}")
try:
response = requests.post(url, headers=headers, json=payload, timeout=10)
response.raise_for_status()
result = response.json()
if result.get("code") != 0:
raise Exception(f"预上传错误: {result.get('msg')}")
return {
"upload_id": result["data"]["upload_id"],
"block_size": result["data"]["block_size"],
"block_num": result["data"]["block_num"]
}
except Exception as e:
raise Exception(f"预上传失败: {str(e)}")
def _create_local_chunks(self, file_path, block_size):
"""在本地生成分片文件"""
# 创建临时目录存储分片
file_name = os.path.basename(file_path)
self.chunk_dir = f"_{file_name}_chunks"
os.makedirs(self.chunk_dir, exist_ok=True)
self._log(f"创建分片目录: {self.chunk_dir}")
chunk_paths = []
with open(file_path, "rb") as f:
seq = 0
while True:
block_data = f.read(block_size)
if not block_data:
break # 读取完毕
# 保存分片到本地
chunk_path = os.path.join(self.chunk_dir, f"chunk_{seq}.part")
with open(chunk_path, "wb") as chunk_file:
chunk_file.write(block_data)
# 计算校验和并保存(用于调试)
checksum = zlib.adler32(block_data) & 0xFFFFFFFF
with open(f"{chunk_path}.checksum", "w") as cs_file:
cs_file.write(str(checksum))
chunk_paths.append(chunk_path)
self._log(f"生成分片 {seq}: {chunk_path} (大小: {len(block_data)}B, checksum: {checksum})")
seq += 1
return chunk_paths
def _upload_local_chunk(self, upload_type, upload_id, seq, chunk_path):
"""上传本地分片文件"""
# 读取本地分片数据
with open(chunk_path, "rb") as f:
block_data = f.read()
block_size = len(block_data)
# 读取预存的校验和(也可重新计算)
checksum_path = f"{chunk_path}.checksum"
if os.path.exists(checksum_path):
with open(checksum_path, "r") as f:
checksum = f.read().strip()
else:
checksum = str(zlib.adler32(block_data) & 0xFFFFFFFF)
url = self.API_ENDPOINTS[upload_type]["upload_part"]
token = self._get_tenant_token()
multipart_data = MultipartEncoder(
fields=[
("upload_id", upload_id),
("seq", str(seq)),
("size", str(block_size)),
("checksum", checksum),
("file", ("blob", block_data, "application/octet-stream"))
]
)
headers = {
"Authorization": f"Bearer {token}",
"Content-Type": multipart_data.content_type,
"Accept": "application/json"
}
try:
response = requests.post(
url,
headers=headers,
data=multipart_data,
timeout=30
)
response.raise_for_status()
result = response.json()
self._log(f"分片{seq}响应: {result}")
if result.get("code") != 0:
raise Exception(f"分片{seq}错误: {result.get('msg')}")
# 兼容飞书data=null的情况
if result.get("data") is None:
self._log(f"分片{seq}上传成功(忽略data=null)")
return None
if "etag" not in result["data"]:
self._log(f"分片{seq}无etag,继续上传")
return None
return result["data"]["etag"]
except Exception as e:
raise Exception(f"分片{seq}上传失败: {str(e)}")
def _finish_upload(self, upload_type, upload_id, block_num):
url = self.API_ENDPOINTS[upload_type]["finish"]
token = self._get_tenant_token()
headers = {
"Authorization": f"Bearer {token}",
"Content-Type": "application/json; charset=utf-8"
}
payload = {"upload_id": upload_id, "block_num": block_num}
try:
response = requests.post(url, headers=headers, json=payload, timeout=10)
response.raise_for_status()
result = response.json()
if result.get("code") != 0:
raise Exception(f"合并错误: {result.get('msg')}")
if "file_token" not in result["data"]:
raise Exception("合并响应无file_token")
return result["data"]
except Exception as e:
raise Exception(f"合并失败: {str(e)}")
def upload(self, file_path, parent_node, upload_type="file", extra=None,
max_retries=3, keep_chunks=False):
"""
先在本地生成分片,再上传
:param keep_chunks: 是否保留本地分片文件(默认删除)
:return: 上传成功的file_token字符串
"""
if not os.path.exists(file_path):
raise FileNotFoundError(f"文件不存在: {file_path}")
file_name = os.path.basename(file_path)
file_size = os.path.getsize(file_path)
if file_size == 0:
raise ValueError("空文件不支持上传")
try:
# 步骤1: 预上传获取分片策略
prepare_result = self._prepare_upload(
upload_type=upload_type,
file_name=file_name,
parent_node=parent_node,
file_size=file_size,
extra=extra
)
upload_id = prepare_result["upload_id"]
block_size = prepare_result["block_size"]
block_num = prepare_result["block_num"]
self._log(f"预上传成功: {upload_id}, 分块数: {block_num}, 分块大小: {block_size}B")
# 步骤2: 在本地生成分片文件
chunk_paths = self._create_local_chunks(file_path, block_size)
if len(chunk_paths) != block_num:
raise Exception(f"本地分片数量不符(预期: {block_num}, 实际: {len(chunk_paths)})")
# 步骤3: 上传本地分片
for seq, chunk_path in enumerate(chunk_paths):
for retry in range(max_retries):
try:
self._log(f"上传分片 {seq + 1}/{block_num} ({chunk_path})")
self._upload_local_chunk(upload_type, upload_id, seq, chunk_path)
break
except Exception as e:
if retry >= max_retries - 1:
raise Exception(f"分片{seq}最终失败: {str(e)}")
self._log(f"分片{seq}重试 {retry + 1}/{max_retries}")
time.sleep(1.5 ** retry)
# 步骤4: 完成上传
result = self._finish_upload(upload_type, upload_id, block_num)
self._log(f"上传成功!file_token: {result['file_token']}")
# 清理本地分片
if not keep_chunks and self.chunk_dir and os.path.exists(self.chunk_dir):
shutil.rmtree(self.chunk_dir)
self._log(f"已删除分片目录: {self.chunk_dir}")
# 直接返回file_token
return result['file_token']
except Exception as e:
print(f"上传失败: {str(e)}")
# 失败时保留分片以便调试
if self.chunk_dir and os.path.exists(self.chunk_dir):
self._log(f"上传失败,保留分片目录: {self.chunk_dir}")
raise
# 使用示例
def main(args):
# 请替换为实际的APP_ID、APP_SECRET、文件路径和目标节点
APP_ID = args["app_id"]
APP_SECRET = args["app_secret"]
FILE_PATH = args["文件路径"]
PARENT_NODE = args["云文件夹token"]
try:
uploader = FeishuLocalChunkUploader(APP_ID, APP_SECRET, debug=args["debug日志打印"])
# 调用upload方法后直接获取file_token
file_token = uploader.upload(
file_path=FILE_PATH,
parent_node=PARENT_NODE,
upload_type="file",
keep_chunks=False # 调试时设为True保留分片
)
print(f"文件上传成功,file_token: {file_token}")
args["file_token"] = file_token
except Exception as e:
print(f"执行错误: {str(e)}")