thảo luận Cách áp dụng LMAX Disruptor với spring

  • Người tạo chủ đề Người tạo chủ đề Mắm Cá Linh
  • Ngày bắt đầu Ngày bắt đầu
Disruptor dùng cho High Frequency Trading thì còn thấy hợp lý. Và đoạn persistence xuống disk thì nên tự làm để low latency chứ lại dùng Kafka để làm input log thì thấy sai vc.
Kafka nổi tiếng là high latency rồi, vì phải wait tất cả các node trong ISR rồi mới được coi là persisted. Một node chậm là kéo theo hệ thống chậm luôn.
Hơn nữa đoạn write xuống MySQL cũng giảm throughput của nó đi nhiều. Vì bản thân MySQL không thể keepup với xử lý trên memory thuần được. Nó vẫn là bottle-neck.
Với cả mấy bàn toán inventory/promotion thì cần gì low latency đến mức phải <1ms
=> Dùng LMAX cho sang chứ k cần thiết.

Áp dụng smart batching mutex lock thông thường thôi là thừa sức rồi
 
Được my friend. Thoải mái luôn.

via theNEXTvoz for iPhone
RingBuffer poll rồi chứa data trong memory. Vậy nếu server crash thì bạn recover data như thế nào?
Mình nghĩ ra mấy hướng sau:
1. Resend data từ message queue + implement idempotency
2. Replicate RingBuffer qua 1 node khác + realtime replication
3. Persist data asynchronously vào 1 database khác, khi processing với RingBuffer rồi track các leaf sequence khi đang process, tạo lại RingBuffer state khi restart + implement idempotency
 
RingBuffer poll rồi chứa data trong memory. Vậy nếu server crash thì bạn recover data như thế nào?
Mình nghĩ ra mấy hướng sau:
1. Resend data từ message queue + implement idempotency
2. Replicate RingBuffer qua 1 node khác + realtime replication
3. Persist data asynchronously vào 1 database khác, khi processing với RingBuffer rồi track các leaf sequence khi đang process, tạo lại RingBuffer state khi restart + implement idempotency
Dùng 1 thread khác pull dữ liệu rồi để vào queue. Nếu crash thì pull lại từ queue.
 
Dùng 1 thread khác pull dữ liệu rồi để vào queue. Nếu crash thì pull lại từ queue.
Thanks bạn.

1. Vậy nếu server không crash thì lúc nào mình clear cái queue này?
2. Sau khi pull lại, mình trực tiếp consume message hay dùng message này để build lại RingBuffer?
a. Nếu consume trực tiếp thì chắc ok, (fail nữa thì cũng chỉ mất 1 message là cùng).
b. Nếu build lại RingBuffer mà server fail tiếp thì làm thế nào? (Cách này có vẻ hơi phức tạp)
 
Thanks bạn.

1. Vậy nếu server không crash thì lúc nào mình clear cái queue này?
2. Sau khi pull lại, mình trực tiếp consume message hay dùng message này để build lại RingBuffer?
a. Nếu consume trực tiếp thì chắc ok, (fail nữa thì cũng chỉ mất 1 message là cùng).
b. Nếu build lại RingBuffer mà server fail tiếp thì làm thế nào? (Cách này có vẻ hơi phức tạp)
có cách khác đây:
https://martinfowler.com/articles/lmax.html

Họ build nguyên 1 con Business Logic Processor khác. Bình thường có 1 con Primary, các con khác standby, vẫn chạy nhưng output bị bỏ qua. Khi con Primary crash thì promote 1 con standby lên làm Primary rồi chạy tiếp. Mình đoán là con standby sẽ save output vào database. Khi recover thì recover data từ đây (chắc process từ timestamp con Primary cũ bị crashed đến lúc promote standby lên làm Primary), hàng ngày clean up data cũ bằng cron job.
 
Như vậy thì Lmax Disruptor giúp triển khai được 1 pipeline có 3 thread riêng biệt Validate Order, Matching Engine
Hi bác, em cũng mới vọc vạch đọc thử thằng Lmax Disruptor, đọc cái comment của bác thì vẫn có 1 thắc mắc như thế này.
Bản thân cái thằng ring buffer nó sẽ lưu data của cùng 1 kiểu, như bác bảo là Validate sau chuyển cho Matching, từ Matching chuyển cho Other Business Logic, thì em đang muốn hỏi thêm bác là: việc "chuyển" ở đây ý bác là như thế nào, em chỉ mới có 2 ý hiểu, nhưng không cái nào giải thích được cái bác nói là 3 thread song song cả.
  1. Thằng Validate gọi Matching, Matching gọi Other Business Logic, nhưng như thế này thì đâu gọi là 3 thằng này chạy song song được
  2. Thằng Validate làm xong nó lưu kết quả vào 1 cái shared data, thằng Matching khi xử lí đến cái data đã được xử lí bởi Validate thì vào check shared data, cái này thì tuỳ logic, có thể là validate pass thfi mới có trong shared data, hoặc vốn trong shared data nó sẽ có 1 field để mark là pass hay fail validation, tương tự với trường hơp other business logic. Nhưng nếu như vậy thì sẽ vẫn block ở shared data.
Mong bác khai sáng thêm.
Updated: Em vừa nghĩ đến 1 khả năng khác là chính cái data lưu trong ring buffer nó đã là 1 cái shared data rồi, mỗi 1 consumer sẽ update vào những field tương ứng của nó. :/ Tuy nhiên em cũng k rõ cái này có phải best practice không.
 
Sửa lần cuối:
Hi bác, em cũng mới vọc vạch đọc thử thằng Lmax Disruptor, đọc cái comment của bác thì vẫn có 1 thắc mắc như thế này.
Bản thân cái thằng ring buffer nó sẽ lưu data của cùng 1 kiểu, như bác bảo là Validate sau chuyển cho Matching, từ Matching chuyển cho Other Business Logic, thì em đang muốn hỏi thêm bác là: việc "chuyển" ở đây ý bác là như thế nào, em chỉ mới có 2 ý hiểu, nhưng không cái nào giải thích được cái bác nói là 3 thread song song cả.
  1. Thằng Validate gọi Matching, Matching gọi Other Business Logic, nhưng như thế này thì đâu gọi là 3 thằng này chạy song song được
  2. Thằng Validate làm xong nó lưu kết quả vào 1 cái shared data, thằng Matching khi xử lí đến cái data đã được xử lí bởi Validate thì vào check shared data, cái này thì tuỳ logic, có thể là validate pass thfi mới có trong shared data, hoặc vốn trong shared data nó sẽ có 1 field để mark là pass hay fail validation, tương tự với trường hơp other business logic. Nhưng nếu như vậy thì sẽ vẫn block ở shared data.
Mong bác khai sáng thêm.
Updated: Em vừa nghĩ đến 1 khả năng khác là chính cái data lưu trong ring buffer nó đã là 1 cái shared data rồi, mỗi 1 consumer sẽ update vào những field tương ứng của nó. :/ Tuy nhiên em cũng k rõ cái này có phải best practice không.
Bạn đang dùng nó à, k em lâu r k dùng nhưng mà mình chưa hiểu ý fren
 
Hi bác, em cũng mới vọc vạch đọc thử thằng Lmax Disruptor, đọc cái comment của bác thì vẫn có 1 thắc mắc như thế này.
Bản thân cái thằng ring buffer nó sẽ lưu data của cùng 1 kiểu, như bác bảo là Validate sau chuyển cho Matching, từ Matching chuyển cho Other Business Logic, thì em đang muốn hỏi thêm bác là: việc "chuyển" ở đây ý bác là như thế nào, em chỉ mới có 2 ý hiểu, nhưng không cái nào giải thích được cái bác nói là 3 thread song song cả.
  1. Thằng Validate gọi Matching, Matching gọi Other Business Logic, nhưng như thế này thì đâu gọi là 3 thằng này chạy song song được
  2. Thằng Validate làm xong nó lưu kết quả vào 1 cái shared data, thằng Matching khi xử lí đến cái data đã được xử lí bởi Validate thì vào check shared data, cái này thì tuỳ logic, có thể là validate pass thfi mới có trong shared data, hoặc vốn trong shared data nó sẽ có 1 field để mark là pass hay fail validation, tương tự với trường hơp other business logic. Nhưng nếu như vậy thì sẽ vẫn block ở shared data.
Mong bác khai sáng thêm.
Updated: Em vừa nghĩ đến 1 khả năng khác là chính cái data lưu trong ring buffer nó đã là 1 cái shared data rồi, mỗi 1 consumer sẽ update vào những field tương ứng của nó. :/ Tuy nhiên em cũng k rõ cái này có phải best practice không.
Data trong ringbuffer thường sẽ là immutable. Muốn truyền từ bước này sang bước khác thì dùng một ringbuffer khác.
Còn nếu muốn làm kiểu update tại chỗ message trong ringbuffer thì sau mỗi bước complete sẽ tăng một atomic variable để báo cho thread khác vị trí đã xử lý done.
Ví dụ có 2 bước validate input và business logic.
Sau khi thành validate input xử lý được 20 message thì sẽ tăng biến atomic lên 20, thằng business logic sẽ có thể read biến đó ra và chỉ xử lý thêm 20 message đó, xử lý xong lại quay lại check biến atomic.
Trong quá trình business logic đang chạy thì thằng validate có thể xử lý tiếp thêm 30 message nữa.
Kiểu vừa song song lại vừa tuần tự
 
Data trong ringbuffer thường sẽ là immutable. Muốn truyền từ bước này sang bước khác thì dùng một ringbuffer khác.
Còn nếu muốn làm kiểu update tại chỗ message trong ringbuffer thì sau mỗi bước complete sẽ tăng một atomic variable để báo cho thread khác vị trí đã xử lý done.
Ví dụ có 2 bước validate input và business logic.
Sau khi thành validate input xử lý được 20 message thì sẽ tăng biến atomic lên 20, thằng business logic sẽ có thể read biến đó ra và chỉ xử lý thêm 20 message đó, xử lý xong lại quay lại check biến atomic.
Trong quá trình business logic đang chạy thì thằng validate có thể xử lý tiếp thêm 30 message nữa.
Kiểu vừa song song lại vừa tuần tự
Còn kiểu các thread check cái biến atomic đó như nào.

Một cách là cứ for loop vô hạn đọc cái giá trị đó lên.
Nếu rất cần low latency thì có thể làm thế. Thậm chí pin cái Thread vào một CPU core để tránh Thread bị đẩy sang core khác mất hết cache. Nhược điểm là CPU lúc nào cũng 100% dù không có message mới để xử lý.

Một cách nữa là có thể sleep một chút, hoặc giống spinlock là chạy một vài instruction kiểu tốn ít điện rồi mới quay lại check lại biến atomic. Hoặc một kiểu nữa là yield lại thread cho OS cho nó có thể schedule thread khác vào.

Một kiểu nữa là dùng mutex + condition variable cơ bản để wait và notify
 
Bạn đang dùng nó à, k em lâu r k dùng nhưng mà mình chưa hiểu ý fren
Em đang muốn làm rõ thêm về cái flow mà bác martin98 đề cập đối với việc áp dụng lmax trong cái exchange-core.
Bác ấy nói cái flow nó là như thế này
Receive event -> EventQueue -> Validate Order (A) -> Matching Engine (B) -> Other Business Logic (C)
Event Queue là 1 cái ring buffer.
Bác ấy nói là (A) sẽ get từ EventQueue ra event để validate, validate xong thì gọi đến (B), (B) xong thì gọi đến (C)
Thì em đang chưa hiểu là ở đây là, vậy cả 3 thằng (A), (B), (C) đều là 3 consumers consume data trên Event Queue, hay chỉ có mỗi thằng (A) thôi.
  1. Nếu chỉ có mỗi thằng (A) thôi thì tại sao em lại cần phải apply cái ring buffer ở đây làm gì? tại sao không chỉ đơn giản là viết 1 cái function trong đó gọi tuần tự (A), (B), (C) mỗi khi có request tới?
  2. Nếu cả 3 đều là consumers, thì làm sao để handle việc là (A) chạy xong, thì đến thằng (B), (C) nó sẽ ignore những thằng fail validation, bởi 3 thằng đó nó sẽ đều duyệt qua cái EventQueue. Một cách đó là phải có 1 cái shared data để thằng (A) chạy xong, nó update thẳng kết quả của validation vào trong shared data đó, thằng (B) khi xử lí event trong EventQueue, nó vào cái shared data đó check xem là kết quả validation là gì, pass thì nó sẽ tiếp tục xử lí event hiện tại, không thì ignore. tương tự với thằng (C). Lúc này thì câu hỏi đặt ra là shared data này tổ chức như thế nào để tránh việc lock.
  3. Hoặc có thể hiểu 1 cách nữa là cả (A), (B) và (C) đều là consumer, nhưng consume data trên các ring buffer riêng biệt. (A) consume event từ EventQueue, xong thì sẽ đẩy event vào ring buffer mà (B) là consumer, tương tự, (B) xong thì đẩy data vào 1 cái ring buffer mà (C) là consumer. Nhưng nếu là như này thì em cũng có thắc mắc tương tự với trường hợp đầu tiên.
 
Data trong ringbuffer thường sẽ là immutable. Muốn truyền từ bước này sang bước khác thì dùng một ringbuffer khác.
Còn nếu muốn làm kiểu update tại chỗ message trong ringbuffer thì sau mỗi bước complete sẽ tăng một atomic variable để báo cho thread khác vị trí đã xử lý done.
Ví dụ có 2 bước validate input và business logic.
Sau khi thành validate input xử lý được 20 message thì sẽ tăng biến atomic lên 20, thằng business logic sẽ có thể read biến đó ra và chỉ xử lý thêm 20 message đó, xử lý xong lại quay lại check biến atomic.
Trong quá trình business logic đang chạy thì thằng validate có thể xử lý tiếp thêm 30 message nữa.
Kiểu vừa song song lại vừa tuần tự
Nếu như bác đề cập thì mình tạo ra 2 cái ring buffer riêng biệt cho validate input và business logic, vì những cái bác đề cập em thấy nó chính là nguyên lí của thằng Disruptor rồi. Còn nếu vẫn muốn chỉ giữ 1 cái ring buffer mà cả 2 thằng validate input và business logic nó đều chạy trên đó thì sẽ hơi nhập nhằng vì mình phải thêm check coi là cái event hiện tại đang consume là được thêm bởi thằng nào.
Nhưng nếu tạo ra 2 cái ring buffer riêng biệt thì em khá là cấn việc tại sao mình lại cần cái ring buffer thay vì việc 1 function gọi thẳng tuần tự validate -> business logic
 
Nếu như bác đề cập thì mình tạo ra 2 cái ring buffer riêng biệt cho validate input và business logic, vì những cái bác đề cập em thấy nó chính là nguyên lí của thằng Disruptor rồi. Còn nếu vẫn muốn chỉ giữ 1 cái ring buffer mà cả 2 thằng validate input và business logic nó đều chạy trên đó thì sẽ hơi nhập nhằng vì mình phải thêm check coi là cái event hiện tại đang consume là được thêm bởi thằng nào.
Nhưng nếu tạo ra 2 cái ring buffer riêng biệt thì em khá là cấn việc tại sao mình lại cần cái ring buffer thay vì việc 1 function gọi thẳng tuần tự validate -> business logic

Ví dụ business logic process bulk 50 request/second nhưng nhất định phải mất 5 second để query metadata từ chỗ khác, còn validate chạy được 10 request/second. Khi validate call business logic trực tiếp, thì mỗi request mình sẽ mất ít nhất 5s. 50 request hết 250s.

Còn nếu để validate với business riêng rẽ thì validate process 50 request hết 5s, business logic process 5s 1 lần, mỗi lần fetch max 50 request, hết 5s. Cả 2 chạy async/parallel thì chạy mất ít hơn 10s.

Btw, mình nghĩ là sẽ tạo ra các ring buffer khác nhau giữa các process.
 
Nếu như bác đề cập thì mình tạo ra 2 cái ring buffer riêng biệt cho validate input và business logic, vì những cái bác đề cập em thấy nó chính là nguyên lí của thằng Disruptor rồi. Còn nếu vẫn muốn chỉ giữ 1 cái ring buffer mà cả 2 thằng validate input và business logic nó đều chạy trên đó thì sẽ hơi nhập nhằng vì mình phải thêm check coi là cái event hiện tại đang consume là được thêm bởi thằng nào.
Nhưng nếu tạo ra 2 cái ring buffer riêng biệt thì em khá là cấn việc tại sao mình lại cần cái ring buffer thay vì việc 1 function gọi thẳng tuần tự validate -> business logic
Cái ý thứ 2 là lý do chạy 2 thread trên cùng một ring buffer đó.
2 cái đó vẫn chạy song song.
Nhưng cái thứ hai luôn đi sau cái thứ nhất 1 bước.
Trong khi cái thứ 2 đang xử lý cái thứ nhất vẫn chạy xử lý message mới bt.
Vừa có song song vừa có một phần tuần tự là thế.
Tuần tự là một msg thì bước 1 xử lý mới được đến bước 2.
Còn song song thì 2 msg khác nhau có thể vừa xử lý bước 1 cho msg này, bước 2 có msg khác.

Call function thì chỉ có tuần tự thôi
 
Cái ý thứ 2 là lý do chạy 2 thread trên cùng một ring buffer đó.
2 cái đó vẫn chạy song song.
Nhưng cái thứ hai luôn đi sau cái thứ nhất 1 bước.
Trong khi cái thứ 2 đang xử lý cái thứ nhất vẫn chạy xử lý message mới bt.
Vừa có song song vừa có một phần tuần tự là thế.
Tuần tự là một msg thì bước 1 xử lý mới được đến bước 2.
Còn song song thì 2 msg khác nhau có thể vừa xử lý bước 1 cho msg này, bước 2 có msg khác.

Call function thì chỉ có tuần tự thôi
Còn làm sao để tránh chạy song song cùng 1 message là dùng biến atomic cho cái sequence number đó.
Biến atomic nó có trong đó cái memory barrier rồi
 

Thống kê chủ đề

Ngày tạo
Mắm Cá Linh,
Người trả lời cuối
quangtung2912,
Trả lời
37
Lượt xem
5.527
Quay lại
Lên đầu trang