 Chess.com 棋局數(shù)據(jù)按月回填與 Filesystem 部分替換(Backfill with Partial Replace))
使用 dlt 實現(xiàn) Chess.com 棋局數(shù)據(jù)按月回填與 Filesystem 部分替換Backfill with Partial Replace【免費下載鏈接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ?項目地址: https://gitcode.com/GitHub_Trending/dl/dlt本文以dlt官方示例 partial_loading.md 為主線講解如何通過 REST API source 聲明式抽取 Chess.com 用戶棋局數(shù)據(jù)并利用文件系統(tǒng)目標Filesystem destination的目錄布局與load_id機制在按月追加新數(shù)據(jù)的同時刪除舊回填文件實現(xiàn)“部分替換式”回填去重。讀完本文你將掌握rest_api_resources的端點配置、response_actions容錯、FilesystemClient的文件遍歷與刪除以及自定義回填清理函數(shù)的完整實戰(zhàn)方案。示例背景為什么需要“部分替換”式回填在數(shù)據(jù)管道中“回填backfill”通常指補錄歷史時間段的數(shù)據(jù)。Chess.com 公開 REST API 按“年/月”提供用戶棋局數(shù)據(jù)例如player/{username}/games/{year}/{month}。當(dāng)我們按月分批把 2023 年全年的對局加載到本地文件系統(tǒng)時會自然產(chǎn)生兩個問題數(shù)據(jù)去重同一個月的文件可能因重復(fù)運行而重復(fù)寫入需要借助主鍵或文件替換機制避免重復(fù)。舊文件清理當(dāng)某一時間段的數(shù)據(jù)被重新加載后之前生成的文件對應(yīng)舊的 load package仍然殘留在表目錄里必須顯式刪除否則讀取該表時會看到“新老文件疊加”的臟數(shù)據(jù)。dlt內(nèi)建的replace寫策略是針對整個表的原子替換而本示例展示的是一種更精細的“部分替換”每個資源對應(yīng)每個月份各自獨立加載append加載完成后由自定義函數(shù)刪除同表目錄下不屬于當(dāng)前l(fā)oad_id的舊文件。這樣既可以保留其它月份的數(shù)據(jù)又能保證同一月份的文件始終來自最新一次加載。源碼總覽示例完整源碼位于 docs/website/docs/examples/partial_loading.md整體流程如下chess_com_sourcedlt.source裝飾的 source遍歷月份列表為每個月通過rest_api_resources生成獨立資源generate_months生成[start_year, start_month]到[end_year, end_month]的連續(xù)月份delete_old_backfills加載完成后刪除目標表目錄中不匹配當(dāng)前l(fā)oad_id的舊文件load_chess_data裝配 pipeline 并執(zhí)行。聲明式配置 Chess.com REST API 源示例使用dlt內(nèi)置的 REST API source 實現(xiàn)聲明式抽取核心入口為rest_api_resources定義于 dlt/sources/rest_api/init.py。它接收一個RESTAPIConfig字典返回資源列表可被dlt.source直接 yield。dlt.source def chess_com_source( username: str, months: list[dict[str, str]] ) - Iterator[DltResource]: for month in months: year month[year] month_str month[month] # Configure REST API endpoint for the specific month config: RESTAPIConfig { client: { base_url: https://api.chess.com/pub/, # Base URL for Chess.com API }, resources: [ { name: fchess_com_games_{year}_{month_str}, # Unique resource name write_disposition: append, endpoint: { path: fplayer/{username}/games/{year}/{month_str}, # API endpoint path response_actions: [ {status_code: 404, action: ignore}, ], }, primary_key: [url], # Primary key to prevent duplicates } ], } yield from rest_api_resources(config)配置要點拆解配置項值作用client.base_urlhttps://api.chess.com/pub/REST API 公共端點前綴無需鑒權(quán)resources[].namechess_com_games_{year}_{month}資源名即目標表名按月唯一resources[].write_dispositionappend每次加載只追加不觸碰其它月份數(shù)據(jù)resources[].endpoint.pathplayer/{username}/games/{year}/{month}按月拉取對局的路徑模板resources[].endpoint.response_actions[{status_code: 404, action: ignore}]無對局的月份404靜默跳過不中斷管道resources[].primary_key[url]以棋局 URL 為主鍵在加載階段做去重值得展開說明的是response_actions的容錯語義。在 dlt/sources/rest_api/config_setup.py 的create_response_hooks實現(xiàn)中所有配置的 response action 會被轉(zhuǎn)換為 HTTP 響應(yīng)鉤子并且默認追加一個raise_for_status鉤子凡是未被response_actions顯式處理的非 2xx 狀態(tài)碼都會拋出 HTTP 錯誤終止該資源。示例中對 404 配置action: ignore對應(yīng)的內(nèi)部行為是拋出IgnoreResponseException從而將該月無數(shù)據(jù)視為正常情況——這對“用戶某月沒下棋”的場景至關(guān)重要。primary_key: [url]則讓 Chess.com 返回的每局棋的url字段成為記錄指紋。dlt在 normalize 階段會依據(jù)主鍵對同批次數(shù)據(jù)做合并去重避免同一局棋因 API 返回重復(fù)或重跑而重復(fù)入庫。與rest_api_source的區(qū)別REST API source 還有另一個入口rest_api_source它返回的是完整的DltSource對象而本示例使用rest_api_resources返回資源列表再在自定義dlt.source中逐月 yield本質(zhì)上是把“REST 配置驅(qū)動”與“自定義循環(huán)生成”兩種方式組合起來實現(xiàn)按月動態(tài)擴展資源集合。從源碼結(jié)構(gòu)看dlt/sources/rest_api/init.pyrest_api_resources與rest_api_source都經(jīng)由rest_api克隆的工廠邏輯完成配置校驗與資源創(chuàng)建二者配置格式完全一致可按需選用。生成連續(xù)月份列表generate_months是一個純粹的時間迭代輔助函數(shù)返回{year: ..., month: ...}字典迭代器月份統(tǒng)一格式化為兩位字符串如01與 API 路徑模板及表名拼接保持一致def generate_months( start_year: int, start_month: int, end_year: int, end_month: int ) - Iterator[dict[str, str]]: start_date p.datetime(start_year, start_month, 1) end_date p.datetime(end_year, end_month, 1) current_date start_date while current_date end_date: yield {year: str(current_date.year), month: f{current_date.month:02d}} # Move to the next month if current_date.month 12: current_date current_date.replace(yearcurrent_date.year 1, month1) else: current_date current_date.replace(monthcurrent_date.month 1)其中p是from dlt.common import pendulum as p引入的 pendulum 時間庫別名dlt 對 pendulum 的封裝見 dlt/common/pendulum.py示例調(diào)用list(generate_months(2023, 1, 2023, 12))即可生成 2023 全年 12 個月的字典列表。該函數(shù)同樣適用于跨年回填如(2022, 11, 2023, 2)會正確依次產(chǎn)出 11、12、1、2 月。配置 Filesystem 目標與管道裝配示例在代碼中直接使用本地文件系統(tǒng)作為目標_storage目錄dest_ dlt.destinations.filesystem(_storage) pipeline dlt.pipeline( pipeline_namechess_com_data, destinationdest_, dataset_namechess_games )dlt.destinations.filesystem(_storage)創(chuàng)建一個本地目錄_storage作為數(shù)據(jù)落地位置pipeline_namechess_com_data管道名用于隔離狀態(tài)與配置dataset_namechess_games數(shù)據(jù)集名文件將寫入_storage/chess_games/下。Filesystem destination 官方文檔docs/website/docs/dlt-ecosystem/destinations/filesystem.md指出其底層基于 fsspec 抽象文件操作因此同一套代碼可平滑切換到 AWS S3、GCS、Azure Blob 等對象存儲只需把bucket_url改為對應(yīng)協(xié)議如s3://your-bucket并在.dlt/secrets.toml中配置憑證。對象存儲的文件布局由 key 模擬語義上與本地目錄一致因此下文的自定義清理邏輯在云存儲上同樣適用。理解文件布局與 load_idFilesystem 目標默認布局為{table_name}/{load_id}.{file_id}.{ext}詳見 filesystem.md其中table_name表名本示例中即chess_com_games_2023_01這類按月生成的資源名load_id本次 load package 的唯一 ID同一批管道運行產(chǎn)生的所有文件共享同一個load_idfile_id同一表同批次的文件序號ext文件格式j(luò)sonl/parquet等。正是“表目錄 load_id命名的文件”這一布局讓示例的自定義刪除邏輯有了可依賴的規(guī)律凡是路徑中不含當(dāng)前l(fā)oad_id的文件都屬于更早的回填批次應(yīng)當(dāng)清理。你也可以在config.toml或代碼中通過layout參數(shù)定制布局例如加入{YYYY}/{MM}時間分區(qū)但若后續(xù)使用replace寫策略布局中必須包含{table_name}占位符且保持前綴規(guī)則否則dlt無法正確推算需刪除的文件集。核心delete_old_backfills 實現(xiàn)部分替換清理示例的精華在于delete_old_backfills函數(shù)——它在每次pipeline.run之后按表執(zhí)行“保留本批、刪除舊批”def delete_old_backfills(load_info: LoadInfo, p: dlt.Pipeline, table_name: str) - None: # Fetch current load id load_id load_info.loads_ids[0] pattern re.compile(rf{load_id}) # Compile regex pattern for the current load ID # Initialize the filesystem client fs_client: FilesystemClient p._get_destination_clients()[0] # type: ignore # Construct the table directory path table_dir os.path.join(fs_client.dataset_path, table_name) # Check if the table directory exists if fs_client.fs_client.exists(table_dir): # Traverse the table directory for root, _dirs, files in fs_client.fs_client.walk(table_dir, maxdepthNone): for file in files: # Construct the full file path file_path os.path.join(root, file) # If the file does not match the current load ID, delete it if not pattern.search(file_path): try: fs_client.fs_client.rm( file_path ) # Remove the old backfill file except Exception as e: print(fError deleting file {file_path}: {e})逐段解讀取當(dāng)前 load_idload_info.loads_ids[0]取自pipeline.run返回的LoadInfo類型定義見 dlt/common/pipeline.py。示例末尾assert len(info.loads_ids) 1也驗證了整次運行只有一個 load package。獲取目標客戶端p._get_destination_clients()實現(xiàn)于 dlt/pipeline/pipeline.py返回當(dāng)前管道的 destination client 元組取第一個即FilesystemClient。該客戶端在 dlt/destinations/impl/filesystem/filesystem.py 中定義封裝了fs_clientfsspec 文件系統(tǒng)句柄與dataset_path數(shù)據(jù)集在存儲中的絕對路徑等屬性。構(gòu)造表目錄os.path.join(fs_client.dataset_path, table_name)定位到chess_games數(shù)據(jù)集下某個月份的表目錄如_storage/chess_games/chess_com_games_2023_01/。遍歷并刪除fs_client.fs_client.walk(table_dir, maxdepthNone)遞歸列出目錄下所有文件對每個文件若完整路徑中不包含當(dāng)前l(fā)oad_id的匹配則調(diào)用fs_client.fs_client.rm(file_path)刪除。刪除異常被捕獲并打印保證單個文件失敗不會中斷整個清理流程。邊界條件與注意事項無目錄即跳過exists檢查保證首次運行或尚未生成該表時不會報錯。加載失敗的文件dlt的文件布局把load_id放進文件名因此同一load_id下的失敗/重試文件也會被一并保留或覆蓋示例按“整個 load package 為準”的粒度清理邏輯自洽??漳夸洑埩魟h除文件后目錄本身未清理對象存儲下空目錄無成本本地文件系統(tǒng)下會留下空文件夾通常可接受。依賴內(nèi)部 API_get_destination_clients以下劃線開頭屬于 dlt 的內(nèi)部接口升級 dlt 大版本時需關(guān)注簽名變化示例代碼中用# type: ignore弱化了類型檢查。主流程按月加載并逐表清理def load_chess_data(): # Initialize the dlt pipeline with filesystem destination, here we use local storage dest_ dlt.destinations.filesystem(_storage) pipeline dlt.pipeline( pipeline_namechess_com_data, destinationdest_, dataset_namechess_games ) # Generate the list of months for the desired date range months list(generate_months(2023, 1, 2023, 12)) # Create the source with all specified months source chess_com_source(MagnusCarlsen, months) # Run the pipeline to fetch and load data info pipeline.run(source) # print(info) # After the run, delete old backfills for each table to maintain data consistency for month in months: table_name fchess_com_games_{month[year]}_{month[month]} delete_old_backfills(info, pipeline, table_name) return info info load_chess_data() assert len(info.loads_ids) 1執(zhí)行過程pipeline.run(source)一次性抽取 12 個月的棋局并按各自表名落地到_storage/chess_games/循環(huán) 12 個月逐表調(diào)用delete_old_backfills刪除這些表目錄中不屬于本次load_id的歷史文件斷言確認本次運行僅產(chǎn)生一個 load package隨后info可被上層使用如打印加載統(tǒng)計。由此第二次重跑腳本時會呈現(xiàn)出“部分替換”效果12 個表目錄內(nèi)的文件全部替換為本輪新加載的文件而數(shù)據(jù)集中其它任何未參與本次回填的表例如后續(xù)追加的 2024 年表完全不受影響。驗證與運行前提運行本示例需要已安裝dlt建議同時安裝 filesystem 依賴pip install dlt[filesystem]可訪問外網(wǎng)以請求https://api.chess.com/pub/該公共端點無需認證示例硬編碼用戶名MagnusCarlsen可替換為任意存在對局記錄的 Chess.com 用戶名。驗證方式運行腳本后檢查_storage/chess_games/下各月份表目錄確認每個目錄中僅剩包含當(dāng)前l(fā)oad_id的文件重復(fù)運行一次觀察舊load_id文件被刪除、新load_id文件保留。相關(guān) REST API source 的完整配置手冊見 REST API source 基礎(chǔ)文檔Filesystem 目標的憑證與布局細節(jié)見 Object store filesystem 文檔。小結(jié)本示例展示了 dlt 生態(tài)中一個可復(fù)用的“部分替換回填”模式聲明式抽取借助rest_api_resources與RESTAPIConfig把按月拉取棋局表達為配置而非手寫 HTTP 客戶端按資源獨立控制寫策略每個月份一個資源、統(tǒng)一append避免replace誤傷其它表基于文件布局的精細清理利用{table_name}/{load_id}...默認布局與FilesystemClient的 walk/rm 能力實現(xiàn)“保留本批、刪除舊批”的冪等回填。這一模式同樣適用于任何“按時間分片、需分批重灌”的文件型數(shù)據(jù)管道日志回灌、指標歷史補數(shù)等只需把“月份”換成你的分片維度即可在文件系統(tǒng)或云對象存儲上實現(xiàn)低成本、可控的部分替換加載。【免費下載鏈接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ?項目地址: https://gitcode.com/GitHub_Trending/dl/dlt創(chuàng)作聲明:本文部分內(nèi)容由AI輔助生成(AIGC),僅供參考