Caideyipi opened a new pull request, #18563:
URL: https://github.com/apache/iotdb/pull/18563

   # Pipe 鍙戦€佸畬鎴愭寚鏍囪璁¤鏄?
   ## 1. 缁撹
   
   鏈疄鐜颁繚鐣欎簡 [Apache IoTDB PR 
#18280](https://github.com/apache/iotdb/pull/18280) 
鐨勫畬鏁寸鍒扮灞忛殰鏂规锛屽苟鍦ㄥ綋鍓嶅垎鏀笂琛ュ厖浜?expected DataRegion 瀹屾暣鎬ф鏌ャ€?
   鏂板鎸囨爣锛?
   ```text
   pipe_datanode_completion_ready{name=<pipeName>,creation_time=<creationTime>}
   ```
   
   鎸囨爣鍊煎彧鏈?`0` 鍜?`1`锛?
   - `1`锛氳繖涓?sender DataNode 涓婏紝璇?Pipe 褰撳墠鎵€鏈夐鏈?DataRegion 鐨勬渶鏂?full-FLUSH 
灞忛殰鍧囧凡缁忚繃 source銆乸rocessor銆乻ink锛屽苟鍦?sink 鎴愬姛 ACK 
鍚庢寜搴忔彁浜わ紱鍚屾椂娌℃湁寰呭鐞嗙殑闈炲績璺充簨浠讹紝涔熸病鏈夋娴嬪埌浠诲姟銆乻ource銆乤ssigner銆乧ommitter銆佸紓甯告垨闄嶇骇鐘舵€佸彉鍖栥€?-
 `0`锛氭湭瀹屾垚銆佷笉鏀寔銆佺姸鎬佹湭鐭ャ€佸彂鐢熺珵鎬佹垨鍋ュ悍妫€鏌ュけ璐ャ€俙0` 涓嶅尯鍒嗗叿浣撳師鍥犮€?
   璁捐鍘熷垯鏄?**鍏佽鏆傛椂鍋囬槾鎬э紝浣嗕笉鍏佽鍋囬槼鎬?*銆傚師鏈?`remaining_event_count` 
缁х画浣滀负杩涘害鎸囨爣锛涘畠涓?`0` 鏄畬鎴愮殑蹇呰鏉′欢锛屼絾涓嶈兘鍗曠嫭浣滀负鈥滃凡鍙戦€佸畬鎴愨€濈殑璇佹槑銆?
   ## 2. 涓轰粈涔堜笉鑳藉彧缁?remaining event 鍔犱笂鈥滅瓑寰?flush鈥濆拰鈥滃凡鎹曡幏 TsFile鈥?
   涓€涓緝灏忕殑瀹炵幇鍙互鎶?storage engine 涓瓑寰呭叧闂殑 processor銆佸凡鎹曡幏鐨?TsFile 鍜?Pipe 
闃熷垪鐩稿姞锛屼絾瀹冨彧鑳藉緱鍒拌繎浼肩Н鍘嬮噺锛屾棤娉曞彲闈犺瘉鏄庡彂閫佸畬鎴愩€?
   ### 2.1 鐘舵€佽縼绉诲瓨鍦ㄨ鏁扮┖绐?
   鍚屼竴涓?TsFile 浼氱粡鍘嗭細
   
   ```text
   working/closing processor -> close callback -> assigner -> source queue -> 
processor queue -> sink/in-flight RPC
   ```
   
   
濡傛灉鍒嗗埆璇诲彇杩欎簺缁撴瀯鍐嶇浉鍔狅紝瀵硅薄鍦ㄤ袱涓粨鏋勪箣闂寸Щ鍔ㄦ椂鍙兘鏆傛椂涓嶅睘浜庝换浣曚竴涓凡璇诲彇鐨勫揩鐓э紝浠庤€岀灛鏃跺緱鍒?`0`銆傜粰鏇村瀹瑰櫒鍔犺鏁板彧鑳界缉灏忕獥鍙o紝涓嶈兘浠庢牴鏈笂娑堥櫎璺ㄧ粍浠跺揩鐓х珵鎬併€?
   ### 2.2 闃熷垪涓虹┖涓嶇瓑浜?sink 宸茬‘璁?
   浜嬩欢浠?sink queue 鍙栧嚭鍚庯紝鍒扮洰鏍囩杩斿洖鎴愬姛 ACK 
涔嬪墠锛岄槦鍒楀彲浠ュ凡缁忎负绌恒€傛鏃?`remaining_event_count == 0` 浠嶄笉鑳借瘉鏄庣洰鏍囩宸茬粡鎺ユ敹鎴栧姞杞藉畬鎴愩€?
   ### 2.3 FLUSH 涓庡苟鍙戝啓鍏ュ瓨鍦ㄨ竟鐣岀珵鎬?
   full FLUSH 寮€濮嬪悗锛屽鏋滄湁骞跺彂 insert 鍒涘缓浜嗘柊鐨?working processor锛岃€岃 processor 
娌¤繘鍏ユ湰娆?flush 鐨勬崟鑾烽泦鍚堬紝鍒欐棫 TsFile 鍏ㄩ儴鍙戦€佸畬涔熶笉鑳戒唬琛ㄦ湰杞啓鍏ュ凡瀹屾垚銆傚繀椤讳负姣忔 full 
FLUSH 寤虹珛 token锛屽苟鍦ㄥ苟鍙?insert 鍑虹幇鏃朵娇 token 澶辨晥銆?
   ### 2.4 鐢熷懡鍛ㄦ湡鍙樺寲浼氳鏃х姸鎬佽鎶?
   Pipe task銆乺ealtime source銆乤ssigner 鎴?committer 
琚浛鎹㈠悗锛屾棫瀹炰緥宸叉彁浜ょ殑鐘舵€佷笉鑳界敤浜庤瘉鏄庢柊瀹炰緥瀹屾垚銆備换鍔″垵濮嬪寲澶辫触閫犳垚 DataRegion 
缂哄け鏃讹紝涔熶笉鑳芥妸鈥滄病鏈夋湰鍦?task鈥濆綋鎴愬畬鎴愩€?
   鍥犳锛屽鏋滅洰鏍囧彧鏄€滆繎浼肩Н鍘嬮噺鈥濓紝灏忔敼 `remaining_event_count` 瓒冲锛涘鏋滅洰鏍囨槸纭畾鈥淧ipe 
宸插彂閫佸畬鎴愨€濓紝鍒欏繀椤诲缓绔嬬鍒扮鏈夊簭灞忛殰銆?18280 
鐨勪富瑕佷綋閲忔鏄敤鏉ュ叧闂笂杩板亣闃虫€х獥鍙o紝涓嶈兘鍙繚鐣欏叾涓煇涓€涓鏁扮偣銆?
   ## 3. 瀹屾垚灞忛殰娴佺▼
   
   ```text
   鍋滄骞?join writers
           |
           v
   鎵ц瑕嗙洊鎵€鏈夌浉鍏?DataRegion 鐨?full FLUSH
           |
           +-- 浣挎棫 completion token 澶辨晥锛岃褰曟柊 token
           +-- 鎹曡幏褰撴椂 working + closing TsFileProcessor
           +-- 绛夊緟宸插湪杩愯鐨勬櫘閫?async flush
           +-- 鍏抽棴骞剁瓑寰呮崟鑾风殑 processor 鍙?close callback
           +-- 浠呭湪 token 鏈骞跺彂 insert 澶辨晥鏃跺彂甯?barrier
           |
           v
   assigner锛堢粦瀹?assigner epoch + data generation锛?        |
           v
   realtime source -> processor锛坆arrier 涓嶅厑璁歌鍚炴帀/鏀瑰啓锛?        |
           v
   sink queue锛坆arrier 涓嶅弬涓庢櫘閫?heartbeat 鍚堝苟锛?        |
           v
   sink 鎴愬姛 ACK -> ordered commit -> onCommitted hook
           |
           v
   completion operator 鍙岄噸蹇収鏍¢獙
           |
           v
   pipe_datanode_completion_ready = 1
   ```
   
   鍏抽敭鐐癸細barrier 鎺掑湪鏈 flush 鎹曡幏骞跺彂甯冪殑 TsFile 浜嬩欢涔嬪悗锛涘彧鏈?sink 鎴愬姛澶勭悊 barrier 
骞跺畬鎴?ordered commit锛孌ataRegion 鎵嶄細琚爣璁板畬鎴愩€?
   ## 4. fail-closed 鏉′欢
   
   浠ヤ笅浠讳竴鏉′欢鎴愮珛锛屾寚鏍囬兘杩斿洖 `0`锛?
   - `remaining_event_count` 涓粛鏈夐潪蹇冭烦浜嬩欢锛?- Pipe 涓嶅瓨鍦ㄣ€佷笉鏄?RUNNING USER 
Pipe锛屾垨瀛樺湪 Pipe/runtime/task 寮傚父锛?- source銆乸rocessor銆乻ink 鎴?TsFile load 
strategy 涓嶅湪鏀寔鑼冨洿锛?- historical source 灏氭湭娑堣垂瀹岋紝鎴?realtime source 灏氭湭瀹屾暣鍚姩锛?- 
褰撳墠瀹為檯 DataRegion task 闆嗗悎涓庢牴鎹?PipeMeta銆乴eader 鍜屾湰鍦?StorageEngine 
璁$畻鍑虹殑棰勬湡闆嗗悎涓嶄竴鑷达紱
   - task/source/assigner/committer 瀹炰緥鍙戠敓鏇挎崲锛?- full-FLUSH token 琚苟鍙?insert 
澶辨晥锛?- barrier 鐨?generation 钀藉悗浜庢渶鏂版暟鎹?generation锛?- 
浜嬩欢鍙戝竷銆佸紩鐢ㄨ鏁般€佸叆闃熸垨渚涚粰鍙戠敓澶辫触锛?- hybrid source 姝e湪绛夊緟 TsFile 鎭㈠宸蹭涪寮冪殑 
tablet锛坉egraded锛夛紱
   - 涓ゆ浠诲姟鎷撴墤蹇収鎴栦袱娆$姸鎬佸揩鐓т箣闂村彂鐢熶换浣曞彉鍖栵紱
   - 鎸囨爣璁$畻鏈熼棿鍙戠敓杩愯鏃跺紓甯告垨鑾峰彇浠诲姟璇婚攣瓒呮椂銆?
   浜嬩欢涓㈠け鎴?publication failure 
灞炰簬涓嶅彲鑷姩璇佹槑鎭㈠鐨勬儏鍐碉紝鎸囨爣浼氫繚鎸?`0`锛岄€氬父闇€瑕佸厛淇闂锛屽啀閲嶅惎/閲嶅缓鐩稿簲 Pipe task 
骞堕噸鏂版墽琛屽畬鎴愬崗璁€?
   ## 5. 鏀寔鑼冨洿
   
   褰撳墠鍙湁浠ヤ笅缁勫悎鑳借繑鍥?`1`锛?
   - RUNNING 鐨?USER Pipe锛?- source/extractor锛歚iotdb-extractor` 
鎴?`iotdb-source`锛?- processor锛歚do-nothing-processor`锛?- 
sink/connector锛氬唴缃?IoTDB Thrift connector/sink 鐨?sync銆乤sync銆丼SL 绛夊埆鍚嶏紱
   - TsFile load strategy锛歚sync`锛?- source 宸插惎鍔紝historical 闃舵宸叉秷璐瑰畬鎴愶紱
   - DataRegion/DML 鍙戦€佸畬鎴愩€?
   鐩墠瀹冧笉鑳借瘉鏄?SchemaRegion/DDL 宸插湪鐩爣绔畬鎴愶紝涔熶笉涓鸿嚜瀹氫箟 processor銆佽嚜瀹氫箟 sink 
鎴栧紓姝?TsFile load strategy 鎻愪緵瀹屾垚淇濊瘉銆傝繖浜涚粍鍚堜細淇濆畧杩斿洖 `0`銆?
   ## 6. 姝g‘浣跨敤鍗忚
   
   1. 鍋滄鎵€鏈夊彲鑳藉悜鐩稿叧 DataRegion 鍐欏叆鐨?writer锛屽苟绛夊緟 writer 绾跨▼缁撴潫銆傚綋鍓?generation 
鏄?DataRegion 绾х殑锛涘嵆浣挎槸 Pipe pattern 涔嬪鐨勫苟鍙戝啓鍏ワ紝涔熷彲鑳戒娇鎸囨爣淇濆畧鍦颁繚鎸?`0`銆?2. 
鎵ц鎴愬姛鐨?full `FLUSH`锛岀‘淇濊鐩栬 Pipe 鍙兘鍙戦€佺殑鎵€鏈?DataRegion銆傛帹鑽愭墽琛屽叏灞€ full 
FLUSH锛涗笉瑕佷娇鐢ㄥ彧鍏抽棴 `SEQ` 鎴?`UNSEQ` 鐨勯儴鍒?flush 浣滀负瀹屾垚灞忛殰銆?3. FLUSH 
鎴愬姛杩斿洖鍚庯紝鍦ㄦ瘡涓鏈?sender DataNode 涓婅幏鍙栨柊椴滅殑鎸囨爣鏍锋湰銆?4. 鍙湁鎵€鏈夐鏈?series 
閮藉瓨鍦ㄤ笖鍊煎潎涓?`1`锛屾墠鑳藉垽瀹氭湰杞?Pipe DML 鍙戦€佸畬鎴愩€?
   Prometheus 鍦烘櫙杩樺繀椤诲悓鏃剁‘璁わ細
   
   - sample timestamp 鏅氫簬鏈 FLUSH锛?- 姣忎釜鐩爣鐨?`up == 1`锛?- series 
鏁伴噺涓庨鏈?sender DataNode 鏁颁竴鑷达紱
   - 涓嶆妸缂哄け series銆佹姄鍙栧け璐ユ垨鏃ф牱鏈綋浣滃畬鎴愩€?
   `creation_time` 鐢ㄤ簬鍖哄垎鍚屽悕 Pipe 鐨勪笉鍚?incarnation銆傚垽鏂椂搴旈攣瀹氭湰杞?Pipe 鐨勫噯纭?`name 
+ creation_time`锛屼笉鑳藉彧鎸?`name` 鑱氬悎鏃?series銆?
   ## 7. 涓?PR #18280 鐨勫叧绯诲強褰撳墠鍒嗘敮閫傞厤
   
   鏈疄鐜伴噰鐢?#18280 鐨勬牳蹇冩彁浜?`12bc3f3a5c00265f5c04dc28a4e6c76affe326bf`锛坄[Pipe] 
Add reliable DataNode completion metric`锛夈€傚畬鏁村疄鐜板寘鍚細
   
   - full FLUSH 瀵?working銆乧losing processor 鍜岄噸鍙?ordinary async flush 鐨勭瓑寰咃紱
   - completion token銆乨ata generation銆乤ssigner epoch 鍜?publication failure 
epoch锛?- completion barrier 鍦?source銆乸rocessor銆乻ink queue 涓殑淇濆簭锛?- sink ACK 
鍚庣殑 ordered commit hook锛?- task/source/assigner/committer 鐢熷懡鍛ㄦ湡鏍¢獙锛?- 
fail-closed completion operator 鍙婂苟鍙戞祴璇曘€?
   褰撳墠鍒嗘敮棰濆澶嶇敤浜?master 宸叉湁鐨?expected-DataRegion 璁$畻锛屽苟瑕佹眰锛?
   ```text
   actual DataRegion source ids == expected DataRegion ids
   ```
   
   杩欓伩鍏嶄簡鏌愪釜 DataRegion task 鍒濆鍖栧け璐ユ垨缂哄け鏃讹紝鎸囨爣鍥犫€滄湰鍦伴泦鍚堜负绌?涓嶅畬鏁粹€濊€岃鎶?`1`銆?
   ## 8. 楠岃瘉缁撴灉
   
   - Spotless锛氶€氳繃锛?- 鐩稿叧妯″潡 clean test-compile锛氶€氳繃锛?- 灞忛殰銆佺敓鍛藉懆鏈熴€侀槦鍒楀拰 
full-flush 骞跺彂閽堝鎬ф祴璇曪細18 涓紝Failures 0锛孍rrors 0锛孲kipped 0锛?- 鑻辨枃 locale 
鍏?reactor `test-compile`锛?2/52 閫氳繃锛?- 涓枃 locale 鍏?reactor `test-compile`锛?2/52 
閫氳繃锛?- `git diff --check`锛氶€氳繃銆?
   閽堝鎬ф祴璇曡鐩栦簡 completion generation銆乫ail-closed銆佹垚鍛樺彉鍖栥€乻ource 鏇挎崲銆乧ommitter 
鏇挎崲銆乥arrier 涓嶈 heartbeat 鍚堝苟銆佷簨浠舵敹闆嗗け璐ャ€乮nsert 浣?flush token 澶辨晥锛屼互鍙?full 
FLUSH 涓?ordinary async flush 閲嶅彔绛夊満鏅€?


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to