RE:NODE

ბაზები11 წუთის საკითხავი

Valkey pub/sub თუ streams: შეტყობინებები და consumer group-ები

როდის გამოიყენო Valkey-ს pub/sub და როდის streams: მიწოდების გარანტიები, consumer group-ები, დადასტურებები, pending ჩანაწერები, trimming და მარცხის შემთხვევები.

0 მკითხველი

Valkey-ს პროცესებს შორის შეტყობინებების გადაცემის ორი გზა აქვს, და ისინი საპირისპირო დაპირებებს იძლევიან. pub/sub არის "გაგზავნე და დაივიწყე": შეტყობინება მიდის იმასთან, ვინც ამ წამს არის გამოწერილი, შემდეგ კი ქრება - შენახვის, დადასტურებისა და ხელახლა გაშვების გარეშე. streams არის მხოლოდ დამატებადი ლოგი: შეტყობინებები ინახება, თითოეულს აქვს id, მომხმარებლებს შეუძლიათ წაკითხვა ნებისმიერი წერტილიდან, ხოლო consumer group-ები აღრიცხავენ, რომელი შეტყობინება დაამუშავა თითოეულმა worker-მა და რომელი არ დაუდასტურებია.

აქედან გამომდინარე წესი მოკლეა. გამოიყენე pub/sub, როცა გამოტოვებულ შეტყობინებას მნიშვნელობა არ აქვს, რადგან შემდეგი მას ცვლის - ცოცხალი შეტყობინებები, ქეშის ინვალიდაციის მინიშნებები, ონლაინ სტატუსი. გამოიყენე streams, როცა ყოველი შეტყობინება ერთხელ მაინც უნდა დამუშავდეს - შეკვეთები, მოვლენები, რომლებიც სხვა სისტემებს ანახლებს, ყველაფერი, რისი დაკარგვაც გაგაბრაზებდა. ორივესთან დაკავშირებული ხარვეზების უმეტესობა მოდის pub/sub-ის გამოყენებიდან იქ, სადაც stream იყო საჭირო. ეს პოსტი ორივეს სათანადოდ ფარავს: ბრძანებებს, გარანტიებს, consumer group-ებს, pending სიას, trimming-ს და მარცხის შემთხვევებს.

განსხვავება ერთ ცხრილში#

Pub/subStreams
ინახებაარაკი, trimming-მდე
მიწოდებამაქსიმუმ ერთხელ, მიმდინარე გამომწერებსერთხელ მაინც, consumer group-ებით
გვიან გამომწერი ძველ შეტყობინებებს ხედავსარაკი, ნებისმიერი id-დან
დადასტურებაარ არისXACK თითო შეტყობინებაზე
დატვირთვის განაწილება worker-ებს შორისარა, ყოველი გამომწერი ყოველ შეტყობინებას იღებსკი, consumer group-ის შიგნით
გადაურჩება Valkey-ს restart-სგადასარჩენი არაფერიაკი, შენახვით
მეხსიერებამხოლოდ ბუფერები ნელი გამომწერებისთვისიზრდება შენახული ჩანაწერებით
ტიპური გამოყენებაცოცხალი განახლებები, ინვალიდაცია, ჩატის გავრცელებამოვლენების ლოგები, სამუშაოს განაწილება, აუდიტის კვალი

არის მესამე ვარიანტი, რომელიც ხალხს ავიწყდება: უბრალო list LPUSH-ითა და BRPOP-ით, რასაც მარტივი ამოცანების რიგები იყენებს. ის ყოველ შეტყობინებას ერთ მომხმარებელს აწვდის და ინახავს, სანამ არ აიღებენ, მაგრამ როგორც კი მომხმარებელმა შეტყობინება ამოიღო, ჩავარდნა მას კარგავს. streams consumer group-ებით სწორედ ამას ასწორებს, და სწორედ ამიტომ არსებობს. ამ პრიმიტივებზე აგებული სრული რიგის ბიბლიოთეკებისთვის იხილე ამოცანების რიგები Valkey-ით.

Pub/sub: ბრძანებები და ქცევა#

გამომქვეყნებელი არხზე აგზავნის; ყოველი კავშირი, რომელიც ამ არხზე ამჟამად გამოწერილია, მას იღებს.

code
# connection ASUBSCRIBE orders:updatesPSUBSCRIBE chat:room:*# connection BPUBLISH orders:updates '{"id":512,"status":"shipped"}'(integer) 1

PUBLISH აბრუნებს გამომწერების რაოდენობას, რომლებმაც შეტყობინება მიიღეს. 0-ის დაბრუნება ნიშნავს, რომ შეტყობინება არსად წავიდა - არავინ უსმენდა, და ის გაქრა. PSUBSCRIBE glob ნიმუშზე იწერს; სასარგებლოა, მაგრამ ყოველი გამოქვეყნებული შეტყობინება ყოველ ნიმუშის გამოწერას ედარება, ამიტომ ათასობით ნიმუში ყოველ გამოქვეყნებაზე CPU-ს ხარჯავს.

რამ, რაც ხალხს აკვირვებს:

  • გამოწერილი კავშირი განსაკუთრებულია. RESP2 პროტოკოლით კავშირს subscribe რეჟიმში მხოლოდ გამოწერასთან დაკავშირებული ბრძანებებისა და PING-ის გაცემა შეუძლია. ამიტომ კლიენტის ბიბლიოთეკები გამოწერისთვის ცალკე კავშირს იყენებენ. ნუ ცდი შენი მთავარი კლიენტის გაზიარებას.
  • გათიშული ნიშნავს გამოტოვებულს. გამომწერი, რომელიც ქსელის შეფერხების, deploy-ის ან garbage collection-ის პაუზის შემდეგ ხელახლა უკავშირდება, ვერაფერს იღებს, რაც მის არყოფნაში გამოქვეყნდა. დაწევა არ არსებობს.
  • ნელი გამომწერები ითიშება. გამომწერის მომლოდინე შეტყობინებები Valkey-ს მეხსიერებაში ბუფერდება. client-output-buffer-limit pub/sub კლიენტებისთვის ნაგულისხმევად 32mb 8mb 60-ია: გამომწერი, რომლის ბუფერი 32 MB-ს აჭარბებს, ან 60 წამის განმავლობაში 8 MB-ზე მაღლა რჩება, ითიშება. ეს სერვერს ერთი გაჭედილი მომხმარებლისგან იცავს, და ნიშნავს, რომ ნელი გამომწერი შეტყობინებებს კარგავს.
  • pub/sub არ ინახება, ამიტომ შენახვა მისთვის არაფერს აკეთებს. AOF და snapshot-ები გასაღებებს ინახავს; არხები გასაღებები არ არის.

sharded pub/sub (SSUBSCRIBE, SPUBLISH) არსებობს cluster-ული deploy-ებისთვის, სადაც ის შეტყობინებებს ერთი shard-ის ფარგლებში ინახავს. ერთ სერვერზე ის ჩვეულებრივი pub/sub-ივით იქცევა.

რისთვის არის pub/sub კარგი#

ამ თვისებების გათვალისწინებით, pub/sub ერგება სამუშაოს, სადაც შეტყობინების დაკარგვა უვნებელია:

  • რეალურ დროში გავრცელება ბრაუზერებზე. რამდენიმე ვებ პროცესი თითოეული WebSocket კავშირებს იჭერს; ერთხელ გამოქვეყნებული მოვლენა ყოველ პროცესს აღწევს, რომელიც მას თავის დაკავშირებულ მომხმარებლებს გადასცემს. Socket.IO-ს Redis adapter და ASP.NET Core SignalR-ის Redis backplane ორივე ასე მუშაობს - SignalR რეალური დროის აპები და WebSocket-ები reverse proxy-ის უკან ამ კონფიგურაციებს ფარავს. თუ მომხმარებელი "ახალი შეტყობინების" სიგნალს გამოტოვებს, შეტყობინებას შემდეგ განახლებაზე დაინახავს.
  • ქეშის ინვალიდაციის მინიშნებები. პროცესები მეხსიერების შიდა ქეშებით cache:invalidate-ზე იწერენ; ჩანაწერი აქვეყნებს შეცვლილ გასაღებს. დააკავშირე ეს ლოკალური ქეშის TTL-თან, რომ გამოტოვებული მინიშნება მხოლოდ ხანმოკლე მოძველებულ მონაცემს ნიშნავდეს.
  • ონლაინ სტატუსი და ბეჭდვის ინდიკატორები. განმარტებით მხოლოდ მათთვისაა აქტუალური, ვინც ახლა უსმენს.
  • დეველოპმენტი და debugging. SUBSCRIBE valkey-cli-ში, რომ მოვლენების ნაკადს უყურო.
fanout.js
import { createClient } from "redis";const pub = createClient({ url: process.env.VALKEY_URL });const sub = pub.duplicate();await Promise.all([pub.connect(), sub.connect()]);await sub.subscribe("orders:updates", (message) => {  broadcastToBrowsers(JSON.parse(message));});await pub.publish("orders:updates", JSON.stringify({ id: 512, status: "shipped" }));

შეამჩნიე duplicate(): მეორე კავშირი გამომწერისთვის, ზემოთ ნახსენები მიზეზით.

Keyspace notifications

Valkey-ს ასევე შეუძლია საკუთარი გასაღებების შესახებ მოვლენების გამოქვეყნება: გასაღების ვადის გასვლა, წაშლა, ჩაწერა. ეს არის keyspace notifications, რომლებიც ჩვეულებრივი pub/sub არხებით მიეწოდება, მაგალითად __keyevent@0__:expired. ნაგულისხმევად ისინი გამორთულია, რადგან მათი გენერირება ყოველ ჩანაწერზე CPU-ს ხარჯავს, და ირთვება notify-keyspace-events პარამეტრით - მაგალითად Ex ვადის გასვლის მოვლენებისთვის.

ისინი მაცდურია "გააკეთე რამე, როცა ამ სესიას ვადა გაუვა" ან "გაუშვი ეს ამოცანა, როცა ტაიმერის გასაღებს ვადა გაუვა" შემთხვევებისთვის, და pub/sub-ის ყოველ სისუსტეს იმემკვიდრეობენ. მსმენელი, რომელიც მოვლენის მომენტში გათიშულია, მის შესახებ ვერასოდეს გაიგებს. ვადის გასვლის მოვლენები ასევე მაშინ ჩნდება, როცა Valkey გასაღებს ნამდვილად შლის, რაც ხდება, როცა გასაღებს მიმართავენ ან მას ვადის ფონური ციკლი პოულობს - და არა აუცილებლად ზუსტად იმ წამს, როცა TTL ამოიწურა. ყველაფრისთვის, რაც აუცილებლად უნდა მოხდეს, იმავე იდეის სანდო ვერსიაა sorted set, რომლის ქულა შესრულების დროა და რომელსაც worker პერიოდულად ამოწმებს. ჰოსტინგზე განთავსებულ ინსტანციაზე notify-keyspace-events-ის შეცვლას CONFIG SET-ის უფლება სჭირდება, რომელიც შენს მომხმარებელს შეიძლება არ ჰქონდეს.

Streams: მხოლოდ დამატებადი ლოგი#

stream არის გასაღები, რომელიც ჩანაწერების დალაგებულ თანმიმდევრობას ინახავს. ყოველი ჩანაწერი ველისა და მნიშვნელობის წყვილების პატარა ნაკრებია, <milliseconds>-<sequence> ფორმის id-ით, რომელსაც სერვერი ანიჭებს.

code
XADD orders * type created id 512 total 49.90"1791468000123-0"XLEN orders(integer) 1XRANGE orders - +1) 1) "1791468000123-0"   2) 1) "type" 2) "created" 3) "id" 4) "512" 5) "total" 6) "49.90"

* Valkey-ს სთხოვს, id თვითონ შექმნას. XRANGE --ითა და +-ით ყველაფერს კითხულობს უძველესიდან უახლესამდე; XREVRANGE უკუღმა კითხულობს. რადგან id-ები დროის ნიშნულებია, შეგიძლია წაიკითხო "ყველაფერი 09:00-დან", თუ საწყისად მილიწამის მნიშვნელობას მისცემ.

consumer group-ების გარეშე ნებისმიერი რაოდენობის მკითხველს შეუძლია stream-ს დამოუკიდებლად მიჰყვეს XREAD-ით, თითოეული იმახსოვრებს ბოლო id-ს, რომელიც ნახა:

code
XREAD COUNT 100 BLOCK 5000 STREAMS orders 1791468000123-0

BLOCK 5000 ახალ ჩანაწერებს ხუთ წამამდე ელოდება; $ id-ის ნაცვლად ნიშნავს "მხოლოდ ჩანაწერები, რომლებიც ამიერიდან დაემატება". ეს უკვე დიდი ნაბიჯია pub/sub-ის შემდეგ: მკითხველს, რომელიც ითიშება, შეუძლია ბოლო id-დან გააგრძელოს და არაფერი გამოტოვოს, სანამ ჩანაწერები trimming-ით არ წაიშლება. ყოველი მკითხველი ყოველ ჩანაწერს ხედავს, ამიტომ ეს გავრცელებაა ხელახლა გაშვების შესაძლებლობით.

ამ ეტაპზე რამდენიმე დიზაინის გადაწყვეტილების მიღება ღირს. ჩანაწერები პატარა და თვითაღწერადი შეინახე: მოვლენის ტიპი, შეცვლილი რამის id და შესაძლოა ვერსია, და არა მთელი ჩანაწერი - მომხმარებლებს მიმდინარე მდგომარეობის ბაზიდან ჩატვირთვა შეუძლიათ, პატარა ჩანაწერები კი stream-ის შენახვას იაფს ტოვებს. გამოიყენე ერთი stream მოვლენის თითო სახეობაზე (orders, payments) და არა ერთი stream ყველაფრისთვის, რომ მომხმარებლებს, რომლებსაც შეკვეთები აინტერესებთ, ყოველი გადახდის წაკითხვა და გამოტოვება არ მოუწიოთ. და id-ები სერვერს მიანიჭებინე; საკუთარის მიწოდება შესაძლებელია, მაგრამ ისინი ყოველთვის უნდა იზრდებოდეს, რისი არევაც რამდენიმე producer-ს შორის მარტივია.

Consumer group-ები: სამუშაოს გაზიარება დადასტურებებით#

ჩანაწერების რამდენიმე worker-ზე გასანაწილებლად ისე, რომ თითოეული ერთხელ დამუშავდეს, შექმენი consumer group. ჯგუფი აღრიცხავს ბოლოს მიწოდებულ id-ს მთლიანად ჯგუფისთვის, პლუს pending entries list-ს (PEL) - შეტყობინებებს, რომლებიც მომხმარებელს მიეწოდა, მაგრამ ჯერ არ დადასტურებულა.

code
XGROUP CREATE orders billing $ MKSTREAMXREADGROUP GROUP billing worker-1 COUNT 10 BLOCK 5000 STREAMS orders >XACK orders billing 1791468000123-0

რას აკეთებს თითოეული ეს სამი ხაზი:

  • XGROUP CREATE ... $ ჯგუფს stream-ის ბოლოდან იწყებს; გამოიყენე 0, რომ დაამუშავო ყველაფერი, რაც მასში უკვე არის. MKSTREAM ქმნის stream-ს, თუ ის ჯერ არ არსებობს.
  • XREADGROUP სპეციალური id-ით > ითხოვს ჩანაწერებს, რომლებიც ჯგუფში არავის მიწოდებია. ყოველი ჩანაწერი ზუსტად ერთ მომხმარებელთან მიდის. მომხმარებლის სახელი (worker-1) თვითნებურია, მაგრამ ერთი და იმავე worker-ის restart-ებს შორის სტაბილური უნდა იყოს.
  • XACK ჩანაწერს pending სიიდან შლის. მანამდე ჯგუფს ახსოვს, რომ ის worker-1-ს აქვს.

რამდენიმე ჯგუფს შეუძლია ერთი და იგივე stream-ის დამოუკიდებლად წაკითხვა: billing და email თითოეული ყოველ შეკვეთას იღებს, ხოლო თითოეული ჯგუფის შიგნით სამუშაო მომხმარებლებს შორის ნაწილდება. ეს არის ნიმუში მოვლენების ლოგისთვის, რომელზეც რამდენიმე ქვესისტემა რეაგირებს.

consumer.py
import os, redisr = redis.Redis.from_url(os.environ["VALKEY_URL"], decode_responses=True)GROUP, NAME = "billing", os.environ.get("WORKER_NAME", "worker-1")try:    r.xgroup_create("orders", GROUP, id="0", mkstream=True)except redis.ResponseError as e:    if "BUSYGROUP" not in str(e):        raisewhile True:    # first, anything delivered to me before a crash and never acknowledged    for stream, entries in r.xreadgroup(GROUP, NAME, {"orders": "0"}, count=50) or []:        for entry_id, fields in entries:            handle(fields); r.xack("orders", GROUP, entry_id)    for stream, entries in r.xreadgroup(GROUP, NAME, {"orders": ">"}, count=50, block=5000) or []:        for entry_id, fields in entries:            handle(fields); r.xack("orders", GROUP, entry_id)

id 0-ით წაკითხვა >-ის ნაცვლად ამ მომხმარებლის საკუთარ pending ჩანაწერებს აბრუნებს - რაც ჩავარდნამდე მიიღო და არასოდეს დაადასტურა. გაშვებისას მათი პირველ რიგში დამუშავება არის გზა, რომლითაც გადატვირთული worker დაწყებულს ასრულებს.

მიწოდების გარანტიები და pending სია#

streams consumer group-ებით ერთხელ მაინც მიწოდებას იძლევა. მომხმარებელი, რომელიც ჩანაწერს ამუშავებს და XACK-მდე ვარდება, მას ხელახლა დაინახავს; ასევე დაინახავს სხვა მომხმარებელი, რომელიც მას დაისაკუთრებს. ამიტომ handler-ები იდემპოტენტური უნდა იყოს - ორჯერ გაშვება უსაფრთხო - ზუსტად ისე, როგორც ნებისმიერ რიგთან.

ჩანაწერები, რომელთა მომხმარებელი სამუდამოდ მოკვდა, სხვამ უნდა დაისაკუთროს:

code
XPENDING orders billingXPENDING orders billing - + 10XAUTOCLAIM orders billing worker-2 60000 0-0 COUNT 10

XPENDING აჯამებს, რა რჩება დაუმუშავებელი, ხოლო დიაპაზონით ჩამოთვლის ჩანაწერებს მათი მფლობელით, უმოქმედობის დროით და მიწოდებების რაოდენობით. XAUTOCLAIM 60,000 ms-ზე მეტხანს უმოქმედო ჩანაწერებს worker-2-ს გადასცემს და აბრუნებს. გაუშვი ის პერიოდულად ყოველ worker-ში, და მკვდარი მომხმარებლებისგან ავტომატური აღდგენა გექნება.

მიწოდებების რაოდენობა შენი მომწამვლელი შეტყობინების დეტექტორია. ჩანაწერი, რომელიც ათჯერ მიეწოდა, მეთერთმეტეზე არ გამოვა; დააკოპირე ის dead-letter stream-ში XADD orders:dead ...-ით, დაადასტურე XACK-ით თავდაპირველ ჯგუფში და გამოსცე გაფრთხილება. ამ ნაბიჯის გარეშე ერთი დაზიანებული შეტყობინება შენს worker-ებში სამუდამოდ ტრიალებს.

მომხმარებლები, რომლებიც აღარ არსებობს, გაასუფთავე XGROUP DELCONSUMER-ით, მას შემდეგ რაც მათ pending ჩანაწერებს დაისაკუთრებ - მომხმარებლის წაშლა მის pending სიას აქრობს.

Trimming, მეხსიერება და შენახვა#

stream ყოველ ჩანაწერს ინახავს, სანამ მას არ წაშლი. ჩანაწერის დადასტურება მას არ შლის - XACK მხოლოდ ჯგუფის pending სიას ანახლებს. trimming-ის გარეშე დატვირთული stream იზრდება, სანამ მეხსიერება არ ამოიწურება.

code
XADD orders MAXLEN ~ 100000 * type created id 513XTRIM orders MAXLEN ~ 100000XTRIM orders MINID ~ 1790863200000

MAXLEN ჩანაწერების რაოდენობას ზღუდავს; MINID შლის ყველაფერს, რაც id-ზე ძველია, რაც დროზე დაფუძნებული შენახვაა ("შეინახე შვიდი დღე"). ~ trimming-ს მიახლოებითს ხდის და Valkey-ს საშუალებას აძლევს, მთელი შიდა კვანძები ერთბაშად წაშალოს; მან შეიძლება ზღვარზე ოდნავ მეტი შეინახოს და ზუსტ trimming-ზე გაცილებით იაფია. გამოიყენე ის, თუ ზუსტი რაოდენობა არ გჭირდება.

trimming არ ამოწმებს, დაამუშავა თუ არა ჩანაწერი ყველა ჯგუფმა. თუ consumer group საკმარისად ჩამორჩება, trimming შლის ჩანაწერებს, რომლებიც მას არასოდეს უნახავს. შენახვის ზომა გულუხვად შეარჩიე იმის მიხედვით, რამდენ ხანს შეიძლება მომხმარებელი გათიშული იყოს, და თვალი ადევნე ჯგუფის ჩამორჩენას - XINFO GROUPS orders აჩვენებს თითოეული ჯგუფის ბოლოს მიწოდებულ id-ს, pending რაოდენობას და, მიმდინარე ვერსიებზე, ჩამორჩენის მაჩვენებელს.

მეხსიერება ჩანაწერზე დამოკიდებულია ველების რაოდენობასა და მნიშვნელობების ზომებზე; პატარა ჩანაწერები თითოეული ბაიტების ათეულები ჯდება, რადგან streams მათ კომპაქტურად ინახავს. ასი ათასი პატარა მოვლენა ჩვეულებრივ რამდენიმე ათეულ მეგაბაიტში ეტევა - გაზომე MEMORY USAGE orders-ით და ნუ ენდობი რაიმე შეფასებას. pub/sub-ისგან განსხვავებით, streams გასაღებებია, ამიტომ AOF და snapshot-ები მათ ინახავს, და eviction პოლიტიკა მათზე ვრცელდება. stream-ს, რომელმაც ჩანაწერები არ უნდა დაკარგოს, noeviction სჭირდება, ზუსტად რიგივით; Valkey-ს შენახვა: AOF და RDB ფარავს დანაკარგის ფანჯარას მკაცრი გაჩერებისას.

RE:NODE-ზე Valkey გეგმა დისკზე ინახავს როგორც AOF-ით, ისე snapshot-ებით, ამიტომ stream-ის ჩანაწერები და consumer group-ის მდგომარეობა restart-ს გადაურჩება. მას პირდაპირ მიმართავ მის ჰოსტსა და პორტზე ამ სერვერისთვის გენერირებული პაროლით; წინ proxy slot არ არის, რაც ზუსტად ისაა, რაც ხანგრძლივ ბლოკირებად წაკითხვებს სჭირდება.

პრობლემების მოგვარება#

გამომწერები შეტყობინებებს ტოვებენ. pub/sub-ით მოსალოდნელია, როცა ისინი გათიშულია, ნელია ან ჯერ არ არის გამოწერილი. თუ ეს მიუღებელია, დიზაინს stream სჭირდება.

`PUBLISH` აბრუნებს 0-ს. იმ მომენტში ამ არხზე გამომწერები არ არის. შეამოწმე არხის სახელი - pub/sub არხები რეგისტრის მიმართ მგრძნობიარეა და ლოგიკური ბაზის ნომრებთან კავშირი არ აქვს.

consumer group-ის ჩამორჩენა სულ იზრდება. მომხმარებლები producer-ებზე ნელია. დაამატე მომხმარებლები ჯგუფს, გაზარდე COUNT, ან იპოვე ნელი handler.

pending სია უსაზღვროდ იზრდება. handler-ები არ იძახებენ XACK-ს, ან მანამდე ვარდებიან. შეამოწმე XPENDING მფლობელებისა და მიწოდებების რაოდენობისთვის.

`BUSYGROUP Consumer Group name already exists`. გაშვებისას უვნებელია; დაიჭირე, როგორც მაგალითში.

მეხსიერება stream-თან ერთად სტაბილურად იზრდება. trimming არ არის. დაამატე MAXLEN ~ XADD-ს, ან გაუშვი XTRIM განრიგით.

FAQ#

შემიძლია მივიღო pub/sub შეტყობინებები, რომლებიც გაიგზავნა, სანამ ჩემი აპი გათიშული იყო?

არა. pub/sub არაფერს ინახავს. თუ ეს გჭირდება, ამის ნაცვლად stream-ში გამოაქვეყნე და ყოველ მკითხველს ბოლო id ააღრიცხვინე; მკითხველები, რომლებიც გათიშული იყვნენ, იქიდან აგრძელებენ, სადაც გაჩერდნენ.

არის Valkey stream Kafka-ს შემცვლელი?

ერთი აპლიკაციისთვის ან რამდენიმე სერვისისთვის ის იმავე ნიმუშს ფარავს - დალაგებულ ლოგს დამოუკიდებელი consumer group-ებით - გაცილებით ნაკლები სამართავით. Kafka აგებულია ტერაბაიტებით გაზომილი შენახვისთვის, ბევრ მანქანაზე გადანაწილებული partition-ებისთვის და კვირების განმავლობაში ხელახლა გაშვებისთვის. თუ შენი stream მეხსიერებაში თავისუფლად ეტევა, Valkey პროპორციული ინსტრუმენტია.

გამოვიყენო streams თუ ამოცანების რიგის ბიბლიოთეკა?

ამოცანების რიგის ბიბლიოთეკა, როცა შეტყობინებები ამოცანებია განმეორებებით, დაყოვნებებითა და დაგეგმვით - BullMQ, Celery და RQ ამას უკვე გაძლევენ. streams, როცა შეტყობინებები მოვლენებია, რომლებზეც რამდენიმე მომხმარებელი რეაგირებს, ან როცა დალაგებული ისტორია გინდა, რომლის ხელახლა გაშვებაც შეგიძლია.

როგორ გავხადო დამუშავება ზუსტად ერთჯერადი?

მხოლოდ transport-ით ვერ გახდი; streams ერთხელ მაინც აწვდის. handler იდემპოტენტური გახადე: ჩაწერე ჩანაწერის id, ან ბიზნეს გასაღები, შენს ბაზაში იმავე ტრანზაქციაში, რაც სამუშაო, და გამოტოვე ყველაფერი, რაც უკვე ჩაწერილია.

ითვლება pub/sub შეტყობინებები მეხსიერების ლიმიტებში?

მხოლოდ სანამ ნელი გამომწერებისთვის ბუფერდება, და ეს ბუფერები თითო კლიენტზე შეზღუდულია client-output-buffer-limit-ით. stream-ის ჩანაწერები სრულად ითვლება, რადგან შენახული მონაცემია.


კომენტარები

სრულიად ანონიმურად: ანგარიშის, ელფოსტის და cookie-ის გარეშე. ინახება მხოლოდ სახელი, ტექსტი და დრო - სხვა არაფერი. ბმულების რაოდენობა ლიმიტირებულია.

0/2000