框架:從定時(shí)任務(wù)到可靠數(shù)據(jù)工作流實(shí)戰(zhàn))
最近在整理幾個(gè)數(shù)據(jù)項(xiàng)目的歷史日志時(shí)我又遇到了那個(gè)熟悉又頭疼的場(chǎng)景每天凌晨系統(tǒng)需要自動(dòng)拉取前一天的交易數(shù)據(jù)進(jìn)行清洗、聚合、計(jì)算幾十個(gè)關(guān)鍵指標(biāo)最后生成報(bào)告。手動(dòng)操作不可能數(shù)據(jù)量太大。寫個(gè)一次性腳本每次跑完就忘下次換個(gè)需求又得重寫。用傳統(tǒng)的定時(shí)任務(wù)工具一旦某個(gè)環(huán)節(jié)出錯(cuò)整個(gè)流程就卡住排查起來像在迷宮里找出口。這其實(shí)就是典型的“批處理作業(yè)”需求。很多開發(fā)者包括早期的我容易陷入一個(gè)誤區(qū)認(rèn)為批處理就是把一堆任務(wù)用cron或systemd timer排個(gè)隊(duì)按時(shí)觸發(fā)就完事了。但真正在生產(chǎn)環(huán)境跑過幾年的人都知道批處理的難點(diǎn)從來不是“啟動(dòng)任務(wù)”而是如何讓一系列任務(wù)可靠、可觀測(cè)、可管理地自動(dòng)運(yùn)行。你需要知道它什么時(shí)候開始、什么時(shí)候結(jié)束、中間每一步是否成功、失敗了怎么重試、資源會(huì)不會(huì)被耗盡、歷史記錄如何追溯。這就是為什么當(dāng)我看到 DolphinDB 的批處理作業(yè)框架時(shí)會(huì)覺得它解決的不是一個(gè)“有沒有”的問題而是一個(gè)“好不好用、穩(wěn)不穩(wěn)定”的問題。它沒有停留在提供一個(gè)簡單的定時(shí)觸發(fā)器而是試圖把數(shù)據(jù)工程師在長期實(shí)踐中積累的那些關(guān)于依賴、容錯(cuò)、監(jiān)控的經(jīng)驗(yàn)沉淀成一套內(nèi)置的、聲明式的系統(tǒng)。今天我們就拋開簡單的“一分鐘學(xué)會(huì)”口號(hào)深入看看這套批處理框架到底在解決什么以及如何把它用對(duì)、用好。1. 批處理作業(yè)的核心從“按時(shí)觸發(fā)”到“可靠完成”很多人對(duì)批處理作業(yè)的第一印象是“定時(shí)跑個(gè)腳本”。這個(gè)理解只對(duì)了一半而且是相對(duì)不重要的一半。定時(shí)只是一個(gè)觸發(fā)條件。批處理作業(yè)真正要管理的是觸發(fā)之后的一連串事件任務(wù)之間的依賴關(guān)系、執(zhí)行過程中的狀態(tài)流轉(zhuǎn)、失敗后的處理策略、以及執(zhí)行歷史的留存與分析。1.1 為什么簡單的定時(shí)任務(wù)不夠用假設(shè)你有一個(gè)經(jīng)典的ETL提取、轉(zhuǎn)換、加載流程從外部API拉取原始數(shù)據(jù)。清洗數(shù)據(jù)處理異常值。將清洗后的數(shù)據(jù)寫入數(shù)據(jù)庫表A?;诒鞟的數(shù)據(jù)進(jìn)行聚合計(jì)算結(jié)果寫入表B。將表B的數(shù)據(jù)導(dǎo)出為CSV報(bào)告。如果你用最基礎(chǔ)的cron來實(shí)現(xiàn)可能會(huì)寫成五個(gè)獨(dú)立的定時(shí)任務(wù)。這立刻會(huì)帶來幾個(gè)問題依賴混亂任務(wù)3必須在任務(wù)2成功完成后才能開始。如果任務(wù)2失敗任務(wù)3卻照常運(yùn)行它會(huì)處理錯(cuò)誤或空數(shù)據(jù)導(dǎo)致后續(xù)結(jié)果全錯(cuò)。你需要在每個(gè)任務(wù)腳本里手動(dòng)檢查上游狀態(tài)代碼迅速變得臃腫。狀態(tài)黑洞任務(wù)跑完了成功還是失敗除了查系統(tǒng)日志可能還很分散沒有集中的視圖。半夜任務(wù)失敗你可能要到第二天早上才發(fā)現(xiàn)。缺乏彈性任務(wù)2因?yàn)榫W(wǎng)絡(luò)波動(dòng)失敗你是希望它立刻重試還是跳過等下次cron本身不提供重試機(jī)制你需要自己在腳本里實(shí)現(xiàn)。資源爭搶如果任務(wù)4非常耗資源而任務(wù)5也同時(shí)被觸發(fā)可能導(dǎo)致系統(tǒng)負(fù)載過高。你需要手動(dòng)錯(cuò)開它們的執(zhí)行時(shí)間或者實(shí)現(xiàn)復(fù)雜的鎖機(jī)制。DolphinDB 的批處理作業(yè)框架本質(zhì)上是在幫你解決這些工程上的“臟活累活”。它讓你能夠以更聲明式的方式描述“要做什么”而把“怎么做”以及“出錯(cuò)怎么辦”交給系統(tǒng)。1.2 DolphinDB 批處理框架的抽象層次DolphinDB 沒有把批處理作業(yè)僅僅看作一個(gè)“任務(wù)”而是將其抽象為一個(gè)有生命周期的對(duì)象。這個(gè)對(duì)象包含幾個(gè)關(guān)鍵維度調(diào)度計(jì)劃不僅僅是“每天幾點(diǎn)”還可以是“每隔N分鐘”、“每周幾”、“每月第幾天”甚至是基于另一個(gè)事件觸發(fā)的復(fù)雜規(guī)則。任務(wù)內(nèi)容具體要執(zhí)行的腳本或函數(shù)。依賴關(guān)系明確指定本任務(wù)需要在哪些其他任務(wù)成功完成后才能啟動(dòng)。重試策略任務(wù)失敗后自動(dòng)重試的次數(shù)、間隔和退避策略。超時(shí)控制防止某個(gè)任務(wù)無限期掛起占用資源。歷史與監(jiān)控每一次執(zhí)行的開始時(shí)間、結(jié)束時(shí)間、狀態(tài)成功/失敗、日志輸出都被系統(tǒng)記錄并提供查詢接口。當(dāng)你用這套框架來描述上面的ETL流程時(shí)你就不再是寫五個(gè)獨(dú)立的cron條目而是定義一個(gè)有向無環(huán)圖DAG。系統(tǒng)會(huì)按照?qǐng)D的依賴關(guān)系有序地推進(jìn)任務(wù)執(zhí)行并自動(dòng)處理狀態(tài)傳遞和故障恢復(fù)。這才是現(xiàn)代批處理作業(yè)該有的樣子。2. 上手第一步超越“Hello World”的最小可行流程官方教程或“一分鐘學(xué)會(huì)”類文章往往從一個(gè)最簡單的定時(shí)打印“Hello World”開始。這有助于理解語法但離真實(shí)場(chǎng)景太遠(yuǎn)容易讓人產(chǎn)生“不過如此”的錯(cuò)覺。我們換個(gè)起點(diǎn)構(gòu)建一個(gè)有實(shí)際意義、有依賴關(guān)系、且能暴露常見問題的最小可行流程。假設(shè)我們有一個(gè)簡化場(chǎng)景每天凌晨計(jì)算前一天的業(yè)務(wù)訂單總額。2.1 環(huán)境準(zhǔn)備與核心對(duì)象認(rèn)知首先確保你的 DolphinDB 服務(wù)已啟動(dòng)并能通過客戶端如 DolphinDB GUI、VS Code 插件或 Python API連接。在 DolphinDB 中批處理作業(yè)的核心是scheduleJob函數(shù)。但直接用它就像直接用底層API比較繁瑣。更常用的方式是使用DailyScheduler或CronScheduler這類更高級(jí)的調(diào)度器對(duì)象。不過為了理解本質(zhì)我們先從基礎(chǔ)入手。一個(gè)完整的作業(yè)定義通常涉及以下幾個(gè)部分作業(yè)函數(shù)封裝具體業(yè)務(wù)邏輯的函數(shù)。調(diào)度器定義何時(shí)觸發(fā)作業(yè)。作業(yè)提交將函數(shù)和調(diào)度器綁定提交給系統(tǒng)。讓我們先創(chuàng)建作業(yè)函數(shù)。這個(gè)函數(shù)需要做幾件事連接數(shù)據(jù)庫如果需要、執(zhí)行查詢、處理結(jié)果、可能還要寫入另一個(gè)表或發(fā)送通知。// 定義作業(yè)函數(shù)計(jì)算昨日訂單總額 def calcYesterdayOrderSum() { // 1. 獲取昨天的日期 yesterday today() - 1 // 2. 假設(shè)我們有一個(gè)訂單表 orderTable包含 orderTime 和 amount 字段 // 這里使用一個(gè)更安全的查詢避免日期邊界問題 sqlQuery select sum(amount) as totalAmount from orderTable where date(orderTime) yesterday // 3. 執(zhí)行查詢 result select * from sqlQuery // 4. 處理結(jié)果這里簡單打印實(shí)際可能寫入結(jié)果表或發(fā)送消息 if (result.size() 0) { total result[totalAmount][0] print([ now() ] 昨日( yesterday )訂單總額為: total) // 可以在這里將 total 寫入一個(gè) daily_summary 表 // ... } else { print([ now() ] 昨日( yesterday )無訂單數(shù)據(jù)。) } }2.2 提交你的第一個(gè)“有狀態(tài)”作業(yè)現(xiàn)在我們使用scheduleJob來提交這個(gè)作業(yè)讓它每天凌晨2點(diǎn)執(zhí)行。// 提交一個(gè)每日定時(shí)作業(yè) jobId scheduleJob(jobIddaily_order_summary, jobDesc計(jì)算昨日訂單總額, jobFunccalcYesterdayOrderSum, scheduleTime02:00m, startDate2024.01.01, endDate2024.12.31) print(作業(yè)已提交ID為: jobId)這里有幾個(gè)關(guān)鍵參數(shù)需要理解jobId: 作業(yè)的唯一標(biāo)識(shí)符必須指定且全局唯一。這是后續(xù)查詢、管理、刪除作業(yè)的依據(jù)。jobDesc: 作業(yè)描述方便人類閱讀。jobFunc: 要執(zhí)行的函數(shù)名。scheduleTime: 每天觸發(fā)的時(shí)間。02:00m表示凌晨2點(diǎn)。startDate/endDate: 作業(yè)的有效期范圍。這個(gè)參數(shù)非常實(shí)用可以用于創(chuàng)建臨時(shí)性的數(shù)據(jù)備份作業(yè)、節(jié)假日特殊處理作業(yè)等。執(zhí)行完上面的代碼作業(yè)就被提交到 DolphinDB 的調(diào)度系統(tǒng)中了。它會(huì)在指定的時(shí)間自動(dòng)觸發(fā)。但這只是開始我們?cè)趺粗浪晒\(yùn)行了2.3 立即驗(yàn)證手動(dòng)觸發(fā)與日志查看不要等到凌晨2點(diǎn)再去驗(yàn)證作業(yè)是否正確。DolphinDB 提供了runJob函數(shù)可以立即手動(dòng)觸發(fā)一個(gè)已提交的作業(yè)。// 立即運(yùn)行指定的作業(yè) runJob(jobId)運(yùn)行后去查看節(jié)點(diǎn)的輸出日志通常在dolphindb.log文件中或者如果你在GUI中在“消息”窗口應(yīng)該能看到我們函數(shù)里print的輸出。這是第一個(gè)實(shí)操建議提交作業(yè)后立刻用runJob手動(dòng)觸發(fā)一次。這能快速驗(yàn)證函數(shù)邏輯是否有語法錯(cuò)誤。函數(shù)是否能訪問到所需的數(shù)據(jù)和表。權(quán)限是否足夠。環(huán)境依賴是否齊全。如果手動(dòng)運(yùn)行都報(bào)錯(cuò)就別指望定時(shí)任務(wù)能成功了。這一步能排除掉80%的初級(jí)問題。3. 從單任務(wù)到工作流構(gòu)建你的第一個(gè)任務(wù)DAG單個(gè)定時(shí)任務(wù)解決了“按時(shí)觸發(fā)”的問題但回到我們開頭的ETL例子真正的挑戰(zhàn)在于任務(wù)間的協(xié)作。下面我們構(gòu)建一個(gè)包含兩個(gè)有依賴關(guān)系的任務(wù)DAG。場(chǎng)景任務(wù)AjobA從模擬數(shù)據(jù)源生成當(dāng)天的訂單明細(xì)并寫入orderTable。任務(wù)BjobB在任務(wù)A成功完成后計(jì)算這些訂單的統(tǒng)計(jì)信息。3.1 定義有依賴關(guān)系的作業(yè)函數(shù)首先定義兩個(gè)作業(yè)函數(shù)。注意jobB需要知道jobA是否成功以及處理的是哪天的數(shù)據(jù)。一種常見的模式是使用“日期分區(qū)”或“狀態(tài)標(biāo)志”。// 作業(yè)A生成當(dāng)日訂單數(shù)據(jù) def generateDailyOrders() { targetDate today() // 生成“今天”的數(shù)據(jù)模擬T1處理 // 模擬生成一些隨機(jī)訂單數(shù)據(jù) n 100 orderTimes datetime(targetDate) rand(86400000, n) // 當(dāng)天隨機(jī)時(shí)間 amounts rand(100.0, n) 50 // 隨機(jī)金額 orderIds “ORD” string(1..n) // 構(gòu)造表并寫入這里假設(shè)orderTable已存在且按日期分區(qū) t table(orderTimes as orderTime, orderIds as orderId, amounts as amount) // 使用append!寫入對(duì)應(yīng)日期的分區(qū) // 注意這里需要根據(jù)你的實(shí)際表結(jié)構(gòu)調(diào)整寫入邏輯 // loadTable(“dfs://orderDB”, “orderTable”).append!(t) print(“[“ now() “] 作業(yè)A已生成” targetDate “日訂單數(shù)據(jù)共” n “條?!? // 關(guān)鍵返回一個(gè)結(jié)果供下游作業(yè)判斷或使用 return targetDate } // 作業(yè)B計(jì)算訂單統(tǒng)計(jì)信息依賴于作業(yè)A的輸出日期 def calcOrderStats(prevJobResult) { // prevJobResult 應(yīng)該是作業(yè)A返回的 targetDate statsDate prevJobResult // 查詢?cè)撊掌诘臄?shù)據(jù)并計(jì)算 // sqlStr select count(*) as cnt, avg(amount) as avgAmt, sum(amount) as totalAmt from loadTable(“dfs://orderDB”, “orderTable”) where date(orderTime) statsDate // result exec cnt, avgAmt, totalAmt from sqlStr // 這里用模擬結(jié)果代替 cnt 100 avgAmt 98.5 totalAmt 9850.0 print(“[“ now() “] 作業(yè)B基于日期” statsDate “計(jì)算統(tǒng)計(jì)訂單數(shù)” cnt “平均金額” avgAmt “總額” totalAmt) // 可以將結(jié)果寫入統(tǒng)計(jì)表 }3.2 使用scheduleJob建立依賴在 DolphinDB 中作業(yè)間的依賴需要通過“前驅(qū)作業(yè)”prevJob參數(shù)來顯式聲明。當(dāng)提交作業(yè)B時(shí)告訴系統(tǒng)它必須在作業(yè)A成功完成后才能運(yùn)行。// 首先提交作業(yè)A每天凌晨1點(diǎn)運(yùn)行 jobAId scheduleJob(jobIdgenerate_orders, jobDesc“生成每日訂單”, jobFuncgenerateDailyOrders, scheduleTime01:00m) // 然后提交作業(yè)B聲明它依賴于 jobA。 // 注意scheduleTime 對(duì)于依賴作業(yè)來說意義變了。它表示在依賴滿足后最早可以開始執(zhí)行的時(shí)間。 // 通常我們會(huì)將其設(shè)置為依賴作業(yè)完成后立即執(zhí)行可以用一個(gè)很早的時(shí)間或者用 after 關(guān)鍵字如果API支持。 // 在DolphinDB當(dāng)前版本更常見的模式是使用 CronScheduler 來組合依賴或者通過判斷上游任務(wù)結(jié)果狀態(tài)表來觸發(fā)。 // 這里演示一種基于完成時(shí)間判斷的思路簡化版 jobBId scheduleJob(jobIdcalc_stats, jobDesc“計(jì)算訂單統(tǒng)計(jì)”, jobFunccalcOrderStats, scheduleTime01:05m, startDate2024.01.01, endDate2024.12.31) print(“作業(yè)B已提交計(jì)劃在每天01:05運(yùn)行但理想情況下應(yīng)在作業(yè)A完成后執(zhí)行。”)這里暴露了一個(gè)關(guān)鍵點(diǎn)原生的scheduleJob在復(fù)雜依賴鏈的表達(dá)上能力有限。它更適合基于固定時(shí)間的調(diào)度。對(duì)于嚴(yán)格的“A成功后再執(zhí)行B”的依賴我們需要更強(qiáng)大的工具——這正是 DolphinDB 的作業(yè)調(diào)度器如DailyScheduler和作業(yè)鏈功能發(fā)力的地方。3.3 邁向工程化使用DailyScheduler管理作業(yè)鏈DailyScheduler提供了更強(qiáng)大的作業(yè)編排能力。我們可以將多個(gè)作業(yè)添加到一個(gè)調(diào)度器中并設(shè)置它們的依賴關(guān)系。// 創(chuàng)建一個(gè)每日調(diào)度器 ds DailyScheduler() // 向調(diào)度器中添加作業(yè)A addJob(ds, jobIdgenerate_orders, jobDesc“生成每日訂單”, jobFuncgenerateDailyOrders, scheduledTime01:00m) // 添加作業(yè)B并指定它必須在 generate_orders 成功后運(yùn)行 addJob(ds, jobIdcalc_stats, jobDesc“計(jì)算訂單統(tǒng)計(jì)”, jobFunccalcOrderStats, scheduledTime01:05m, dependencies[generate_orders]) // 提交整個(gè)調(diào)度器 submit(ds)通過dependencies[generate_orders] 參數(shù)我們清晰地定義了作業(yè)B對(duì)作業(yè)A的依賴。調(diào)度器會(huì)負(fù)責(zé)管理執(zhí)行順序每天凌晨1點(diǎn)嘗試執(zhí)行g(shù)enerate_orders。只有g(shù)enerate_orders成功完成函數(shù)正常返回未拋出異常調(diào)度器才會(huì)在1:05或依賴滿足后立即觸發(fā)calc_stats。如果generate_orders失敗calc_stats將不會(huì)被執(zhí)行。這種方式才真正實(shí)現(xiàn)了我們想要的有向無環(huán)圖DAG工作流。你可以構(gòu)建更復(fù)雜的鏈條比如[A] - [B][A] - [C][B, C] - [D]。4. 保障與洞察讓批處理作業(yè)變得可觀測(cè)、可管理作業(yè)提交并運(yùn)行起來只是萬里長征第一步。在生產(chǎn)環(huán)境中你需要回答以下問題昨晚的批處理跑完了嗎哪個(gè)環(huán)節(jié)失敗了為什么每個(gè)任務(wù)花了多長時(shí)間歷史執(zhí)行記錄能保存多久如何查詢DolphinDB 的批處理框架內(nèi)置了這些運(yùn)維能力的支持。4.1 監(jiān)控作業(yè)執(zhí)行狀態(tài)系統(tǒng)提供了若干函數(shù)來查詢作業(yè)信息// 1. 查看所有已提交的作業(yè)包括一次性作業(yè)和定時(shí)作業(yè) getScheduledJobs() // 2. 查看最近N次的作業(yè)執(zhí)行記錄非常有用 getJobHistory(10) // 查看最近10條執(zhí)行記錄 // 返回的表格通常包含jobId, startTime, endTime, status, message // status 可能是 ‘成功’、‘失敗’、‘運(yùn)行中’ // message 可能包含錯(cuò)誤信息或打印輸出 // 3. 查看特定作業(yè)的下次執(zhí)行時(shí)間 getJobSchedule(daily_order_summary)養(yǎng)成習(xí)慣每天上班第一件事先跑一下getJobHistory(50)??焖贋g覽一下狀態(tài)列是否有“失敗”的記錄。這是最基礎(chǔ)的批處理作業(yè)健康檢查。4.2 處理失敗與實(shí)現(xiàn)重試任務(wù)失敗是常態(tài)。網(wǎng)絡(luò)抖動(dòng)、資源不足、臨時(shí)鎖、數(shù)據(jù)異常都可能導(dǎo)致失敗。一個(gè)健壯的批處理系統(tǒng)必須能處理失敗。在scheduleJob或addJob時(shí)可以配置重試策略// 在 addJob 時(shí)指定重試策略示例具體參數(shù)名請(qǐng)查閱最新版本文檔 addJob(ds, jobIdfetch_external_data, jobFuncfetchData, scheduledTime00:30m, maxRetries3, retryInterval60)參數(shù)解讀maxRetries3最多自動(dòng)重試3次不含首次執(zhí)行。retryInterval60每次重試間隔60秒。重試策略的選擇是一門學(xué)問立即重試適用于因瞬時(shí)鎖、線程競(jìng)爭導(dǎo)致的失敗。間隔可以很短如10秒。延遲重試適用于依賴外部服務(wù)如API暫時(shí)不可用。間隔可以長一些如5分鐘。指數(shù)退避更高級(jí)的策略每次重試間隔時(shí)間指數(shù)級(jí)增加避免對(duì)故障服務(wù)造成“驚群”效應(yīng)。DolphinDB 可能通過其他參數(shù)或自定義函數(shù)支持。重要提醒不是所有失敗都適合重試。如果是業(yè)務(wù)邏輯錯(cuò)誤如SQL語法錯(cuò)誤、數(shù)據(jù)格式永久性錯(cuò)誤重試多少次都會(huì)失敗。這時(shí)作業(yè)會(huì)達(dá)到最大重試次數(shù)后最終失敗并留下錯(cuò)誤日志。你需要根據(jù)getJobHistory中的錯(cuò)誤信息 (message) 進(jìn)行人工排查和修復(fù)。4.3 管理作業(yè)生命周期作業(yè)不是提交了就一勞永逸。業(yè)務(wù)邏輯會(huì)變調(diào)度需求也會(huì)變。// 1. 刪除一個(gè)作業(yè) deleteJob(daily_order_summary) // 2. 暫停一個(gè)作業(yè)使其不再被調(diào)度 pauseJob(daily_order_summary) // 3. 恢復(fù)一個(gè)被暫停的作業(yè) resumeJob(daily_order_summary) // 4. 立即觸發(fā)一次作業(yè)運(yùn)行用于測(cè)試或補(bǔ)數(shù)據(jù) runJob(daily_order_summary) // 5. 修改作業(yè)的調(diào)度時(shí)間或參數(shù)通常需要先刪除再重新提交對(duì)于使用DailyScheduler提交的作業(yè)鏈管理單元是整個(gè)調(diào)度器。你可以暫停、恢復(fù)或刪除整個(gè)調(diào)度器從而控制其中所有作業(yè)。4.4 將作業(yè)日志接入你的監(jiān)控系統(tǒng)生產(chǎn)環(huán)境的運(yùn)維往往需要一個(gè)集中的監(jiān)控平臺(tái)如 Prometheus Grafana。DolphinDB 的作業(yè)執(zhí)行記錄本身存儲(chǔ)在系統(tǒng)表中如JOB_HISTORY你可以定期將這些數(shù)據(jù)導(dǎo)出或者通過 DolphinDB 的 API 被外部系統(tǒng)拉取從而在統(tǒng)一的看板上展示批處理作業(yè)的健康狀態(tài)、執(zhí)行時(shí)長趨勢(shì)等。更進(jìn)階的做法是在作業(yè)函數(shù)中將關(guān)鍵里程碑開始、成功、失敗和性能指標(biāo)耗時(shí)、處理數(shù)據(jù)量寫入一個(gè)專門的監(jiān)控表或發(fā)送到消息隊(duì)列實(shí)現(xiàn)更細(xì)粒度的監(jiān)控和告警。批處理作業(yè)從“能跑”到“跑得穩(wěn)”核心就在于這些運(yùn)維細(xì)節(jié)的打磨。DolphinDB 提供了基礎(chǔ)的工具和框架而如何利用好它們構(gòu)建出適合自己業(yè)務(wù)場(chǎng)景的、可靠的數(shù)據(jù)流水線則需要我們根據(jù)上述原則去設(shè)計(jì)和實(shí)踐。記住好的批處理系統(tǒng)是讓數(shù)據(jù)工程師在晚上能睡個(gè)安穩(wěn)覺的系統(tǒng)。