飞书分片上传实现PY代码分享
评论
收藏

飞书分片上传实现PY代码分享

经验分享
羽鹿
2025-10-16 14:32·浏览量:1217
羽鹿
影刀高级开发者
发布于 2025-10-16 14:21更新于 2025-10-16 14:321217浏览

背景

对于特大文件上传到飞书的功能实现

直接贴源代码,文件 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)}")


收藏
全部评论1
最新
发布评论
评论