Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
Ez a szakasz a Stream implementációjának magas szintű áttekintését nyújtja Orleans . Az alkalmazás szintjén nem látható fogalmakat és részleteket ismerteti. Ha csak streameket szeretne használni, nem kell elolvasnia ezt a szakaszt.
Terminológia:
A "queue" szóra hivatkozunk minden olyan tartós tárolási technológiára, amely képes befogadni a stream eseményeket, és amely lehetővé teszi események lekérését, vagy leküldéses mechanizmust biztosít az események fogyasztásához. A méretezhetőség érdekében ezek a technológiák általában darabolt/particionált üzenetsorokat biztosítanak. Az Azure Queues például lehetővé teszi, hogy több üzenetsort hozzon létre, és az Event Hubs pedig több központot is.
Állandó folyamatok
Minden Orleans állandó streamszolgáltató közös implementációval PersistentStreamProviderrendelkezik. Ezeket az általános streamszolgáltatókat technológiaspecifikus IQueueAdapterFactorykonfigurálással kell konfigurálni.
Tesztelési célokra például olyan üzenetsor-adapterekkel rendelkezünk, amelyek a tesztelési adatokat generálják ahelyett, hogy egy üzenetsorból olvasnánk be az adatokat. Az alábbi kód bemutatja, hogyan konfiguráljuk a perzisztens streamszolgáltatót az egyéni (generátoros) üzenetsor-adapter használatára. Ezt úgy teszi, hogy konfigurálja az állandó streamszolgáltatót az adapter létrehozásához használt gyári függvénnyel.
hostBuilder.AddPersistentStreams(
StreamProviderName, GeneratorAdapterFactory.Create);
Amikor egy stream előállító létrehoz egy új streamelemet, és meghívja a stream.OnNext()-t, a Orleans stream futási környezet végrehajtja a megfelelő metódust a IQueueAdapter stream szolgáltatónak, amely közvetlenül a megfelelő üzenetsorra helyezi az elemet.
Ügynökök lekérése
Az állandó streamszolgáltató középpontjában a lekérési ügynökök állnak. Az eseményhúzó ügynökök tartós üzenetsorokból húzzák ki az eseményeket, és azokat az alkalmazáskód szemcséibe továbbítják, amelyek felhasználják őket. A lekérdező ügynököket elosztott "mikroszolgáltatásnak" tekinthetjük – particionált, magas rendelkezésre állású és rugalmas elosztott összetevőként. A pull ügynökök ugyanabban a silóban futnak, amelyek az alkalmazás moduljait üzemeltetik, és teljes mértékben a Orleans Streaming Runtime által felügyeltek.
StreamQueueMapper és StreamQueueBalancer
A lekéréses ügynökök IStreamQueueMapper és IStreamQueueBalancer paramétereivel vannak ellátva. A IStreamQueueMapper az összes várólista felsorolását tartalmazza, és az adatfolyamok várólistákhoz való hozzárendeléséért is felelős. Így a Persistent Stream Provider gyártói oldala tudja, hogy melyik sorba kell beilleszteni az üzenetet.
A IStreamQueueBalancer azt fejezi ki, hogyan van egyensúlyban a sorok kezelése a Orleans silók és ügynökök között. A cél az, hogy kiegyensúlyozott módon rendeljen várólistákat az ügynökökhöz, hogy megelőzze a szűk keresztmetszeteket és támogassa a rugalmasságot. Amikor új silót ad hozzá a Orleans fürthöz, a rendszer automatikusan újraegyensúlyozza az üzenetsorokat a korábbi és az új silók között. Ez StreamQueueBalancer lehetővé teszi a folyamat testreszabását.
Orleans számos beépített StreamQueueBalancer-sal rendelkezik, a különböző kiegyensúlyozási forgatókönyvek (nagy és kis számú üzenetsor) és a különböző környezetek (Azure, helyszíni, statikus) támogatásához.
A fenti tesztgenerátor példáját használva az alábbi kód bemutatja, hogyan konfigurálható az üzenetsor-leképező és az üzenetsor-kiegyensúlyozó.
hostBuilder
.AddPersistentStreams(StreamProviderName, GeneratorAdapterFactory.Create,
providerConfigurator =>
providerConfigurator.Configure<HashRingStreamQueueMapperOptions>(
ob => ob.Configure(options => options.TotalQueueCount = 8))
.UseDynamicClusterConfigDeploymentBalancer());
A fenti kód konfigurálja a GeneratorAdapterFactory-t egy nyolc üzenetsoros leképező használatára, és egyensúlyba hozza az üzenetsorokat a fürtön belül a DynamicClusterConfigDeploymentBalancer segítségével.
Letöltési protokoll
Minden siló egy sorból húzó ügynököket futtat, minden ügynök egy adott sorból húz. A lekéréseket végrehajtó ügynökök egy belső futtatókörnyezeti összetevő, a SystemTarget által implementálva vannak. A SystemTargetek alapvetően futtatókörnyezeti szemcsék, egyszálas egyidejűségnek vannak kitéve, normál szemcseküldést használhatnak, és ugyanolyan könnyűek, mint a szemcsék. A szemcsékkel ellentétben a SystemTargetek nem virtuálisak: a futtatókörnyezet által explicit módon jönnek létre, és nincsenek helyfüggetlenek. Az adatátviteli ügynökök SystemTargets-ként történő megvalósításával a Orleans adatfolyam-futtatókörnyezet a beépített Orleans funkciókra támaszkodhat, és képes nagy számú sor kezelésére, mivel az új adatátviteli ügynök létrehozása ugyanannyira költséghatékony, mint egy új egység létrehozása.
Minden lekéréses ügynök egy időzített eseményt futtat, amely a IQueueAdapterReceiver.GetQueueMessagesAsync metódus meghívásával lehív az üzenetváró sorból. A visszaadott üzenetek az ügynökönkénti belső adatstruktúrába kerülnek.IQueueCache A rendszer minden üzenetet megvizsgál, hogy megtudja a célstreamjét. Az ügynök a Pub-Sub használja a streamre feliratkozott streamfelhasználók listájának megkeresésére. A fogyasztói lista lekérése után az ügynök helyileg tárolja azt (a pub-sub cache-ben), így nem kell minden egyes üzenetnél konzultálnia Pub-Sub-val. Az ügynök a pub-sub-ra is feliratkozik, hogy értesítést kapjon az új fogyasztókról, akik feliratkoznak az adott streamre. Ez az ügynök és a kiadó-előfizetői modell közötti kézfogás erős stream-előfizetési szemantikát garantál: miután a fogyasztó feliratkozott a streamre, látni fogja az összes eseményt, amely a feliratkozás után lett létrehozva. Emellett a StreamSequenceToken használata lehetővé teszi, hogy valamilyen korábbi időpontra előfizessen.
Üzenetsor-gyorsítótár
IQueueCache Egy ügynökönkénti belső adatstruktúra, amely lehetővé teszi az új események üzenetsorból való leválasztását, és eljuttatja azokat a fogyasztókhoz. Emellett lehetővé teszi a különböző streamek és különböző fogyasztók számára történő kézbesítés szétválasztását is.
Képzeljen el egy olyan helyzetet, amikor egy stream 3 streamfelhasználóval rendelkezik, és az egyik lassú. Ha nem gondoskodnak róla, ez a lassú felhasználó befolyásolhatja az ügynök előrehaladását, lelassíthatja a stream más felhasználóinak használatát, és akár lassíthatja az események lekérdezését és továbbítását más streamek esetében is. Ennek megakadályozása és az ügynök maximális párhuzamosságának lehetővé tétele érdekében használjuk a IQueueCache-t.
IQueueCache puffereli az eseményeket, és lehetővé teszi az ügynök számára, hogy az eseményeket a saját tempójában kézbesítse az egyes fogyasztóknak. A fogyasztónkénti kézbesítést a belső, úgynevezett IQueueCacheCursorösszetevő valósítja meg, amely nyomon követi a fogyasztónkénti előrehaladást. Így minden fogyasztó a saját tempójában fogadja az eseményeket: a gyors fogyasztók olyan gyorsan kapják meg az eseményeket, ahogy kikerülnek a sorból, míg a lassú fogyasztók később fogadják őket. Miután az üzenet az összes felhasználóhoz el lett küldve, törölhető a gyorsítótárból.
Visszanyomás
A streamelési futtatókörnyezetben a Orleans visszanyomás két helyen érvényes: streames események továbbítása az üzenetsorból az ügynökhöz , és az események továbbítása az ügynöktől a streamfelhasználók számára.
Az utóbbit a beépített üzenetkézbesítési Orleans mechanizmus biztosítja. Az ügynök minden streameseményt a standard Orleans szemcsés üzenetküldés segítségével, egyenként továbbít a fogyasztókhoz. Vagyis az ügynökök egy eseményt (vagy korlátozott méretű eseményköteget) küldenek minden streamfelhasználónak, és várják ezt a hívást. A következő esemény csak akkor indul el, ha az előző esemény feladatát feloldották vagy megszakították. Így természetesen a fogyasztónkénti kézbesítési arányt egyszerre egy üzenetre korlátozzuk.
Amikor adatfolyam eseményeket juttat az üzenetsorból az ügynökhöz, Orleans az adatfolyam egy új, speciális visszatartási mechanizmust biztosít. Mivel az ügynök leválasztja az események lekérését az üzenetsorról, és eljuttatja őket a fogyasztóknak, egyetlen lassú fogyasztó olyannyira lemaradhat, hogy a IQueueCache rendszer feltöltődik. A határozatlan ideig történő növekedés megakadályozása IQueueCache érdekében korlátozzuk a méretét (a méretkorlát konfigurálható). Az ügynök azonban soha nem dobja el a nem kézbesített eseményeket.
Ehelyett, amikor a gyorsítótár elkezd töltődni, az ügynökök lelassítják az események lekérésének sebességét az üzenetsorból. Így a lassú kézbesítési időszakokat úgy vészelhetjük át, hogy módosítjuk az üzenetsorból történő felhasználás rátáját (backpressure), majd később visszatérhetünk a gyors fogyasztási arányokhoz. A "lassú kézbesítés" völgyeinek észleléséhez a IQueueCache gyorsítótár-gyűjtők belső adatstruktúráját használja, amely nyomon követi az eseményeknek az egyes streamfelhasználókhoz való továbbításának előrehaladását. Ez egy nagyon rugalmas és önkorrekt rendszert eredményez.