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 referenciaarchitektúra egy végpontok közötti streamfeldolgozási folyamatot mutat be. A folyamat két forrásból betölti az adatokat, korrelálja a két stream rekordjait, és kiszámítja a gördülő átlagot egy időablakban. Az eredmények tárolása további elemzés céljából történik.
Architektúra
Töltse le az architektúra Visio-fájlját.
Workflow
Az architektúra az alábbi összetevőkből áll:
Adatforrások. Ebben az architektúrában két adatforrás létezik, amelyek valós időben hoznak létre adatfolyamokat. Az első stream tartalmazza a menetadatokat, a második pedig a viteldíjadatokat. A referenciaarchitektúra egy szimulált adatgenerátort tartalmaz, amely statikus fájlokból olvas be, és leküldi az adatokat az Event Hubsba. Egy valós alkalmazásban az adatforrások a taxifülkékbe telepített eszközök.
Azure Event Hubs. Az Event Hubs egy eseménybetöltési szolgáltatás. Ez az architektúra két eseményközpont-példányt használ, egyet minden adatforráshoz. Minden adatforrás adatstreamet küld a társított eseményközpontnak.
Azure Stream Analytics. A Stream Analytics egy eseményfeldolgozó motor. A Stream Analytics-feladat beolvassa az adatfolyamokat a két eseményközpontból, és elvégzi a streamfeldolgozást.
Azure Cosmos DB. A Stream Analytics-feladat kimenete rekordsorozat, amely JSON-dokumentumként íródik egy Azure Cosmos DB-dokumentum-adatbázisba.
Microsoft Power BI. A Power BI üzleti elemzési eszközökből álló csomagja az üzleti elemzésekhez szükséges adatok elemzéséhez. Ebben az architektúrában betölti az Adatokat az Azure Cosmos DB-ből. Így a felhasználók elemezhetik az összegyűjtött előzményadatok teljes készletét. Az eredményeket közvetlenül streamelheti a Stream Analyticsből a Power BI-ba az adatok valós idejű megtekintéséhez. További információ: Valós idejű streamelés a Power BI-ban.
Azure Monitor. Az Azure Monitor a megoldásban üzembe helyezett Azure-szolgáltatások teljesítménymetrikáit gyűjti. Ha ezeket egy irányítópulton vizualizálja, betekintést nyerhet a megoldás állapotába.
Forgatókönyv részletei
Forgatókönyv: A taxitársaság minden taxiútról adatokat gyűjt. Ebben a forgatókönyvben feltételezzük, hogy két különálló eszköz küld adatokat. A taxinak van egy mérőórája, amely információkat küld az egyes utazásokról — az időtartamról, a távolságról, valamint a felszállási és leszállási helyszínekről. Egy külön eszköz fogadja az ügyfelektől érkező kifizetéseket, és adatokat küld a viteldíjakról. A taxis társaság szeretné valós időben kiszámítani az egy mérföldre eső átlagos borravalót, hogy felismerje a trendeket.
Lehetséges használati esetek
Ez a megoldás a kiskereskedelmi forgatókönyvhöz van optimalizálva.
Adatok betöltése
Az adatforrás szimulálásához ez a referenciaarchitektúra a New York-i taxiadatkészletet[1] használja. Ez az adatkészlet a New York-i taxiutakra vonatkozó adatokat tartalmazza négyéves időszak alatt (2010–2013). Kétféle rekordot tartalmaz: a menetadatokat és a viteldíjakat. A menetadatok magukban foglalják az utazás időtartamát, a menet távolságát, valamint a felvételi és a kiszállási helyet. A viteldíjadatok tartalmazzák a viteldíjakat, az adó- és a tippösszegeket. Mindkét rekordtípus gyakori mezői közé tartozik a medálszám, a feltört licenc és a szállító azonosítója. Ez a három mező együttesen azonosítja a taxit és a sofőrt. Az adatok CSV formátumban lesznek tárolva.
[1] Donovan, Brian; Work, Dan (2016): New York-i taxiutazási adatok (2010-2013). Illinois-i Egyetem, Urbana-Champaign. https://doi.org/10.13012/J8PN93H8
Az adatgenerátor egy .NET alkalmazás, amely beolvassa a rekordokat, és elküldi őket Azure Event Hubs. A generátor JSON formátumban küldi el az utazási adatokat, a menetdíjakat pedig CSV formátumban.
Az Event Hubs partíciókkal szegmentálta az adatokat. A partíciók lehetővé teszik, hogy a felhasználó párhuzamosan olvassa be az egyes partíciókat. Amikor adatokat küld az Event Hubsnak, explicit módon megadhatja a partíciókulcsot. Ellenkező esetben a rekordok körbe haladó módon vannak hozzárendelve a partíciókhoz.
Ebben a konkrét forgatókönyvben a menetadatoknak és a viteldíjadatoknak egy adott taxihoz ugyanazzal a partícióazonosítóval kell rendelkezniük. Ez lehetővé teszi, hogy a Stream Analytics bizonyos fokú párhuzamosságot alkalmazzon a két stream korrelációja esetén. A menetadatok n partíciójában lévő rekord megegyezik a viteldíjadatok n partíciójában lévő rekorddal.
Az adatgenerátorban a két rekordtípus közös adatmodellje rendelkezik olyan PartitionKey tulajdonsággal, amely a Medallion, HackLicense és VendorId összefűzése.
public abstract class TaxiData
{
public TaxiData()
{
}
[JsonProperty]
public long Medallion { get; set; }
[JsonProperty]
public long HackLicense { get; set; }
[JsonProperty]
public string VendorId { get; set; }
[JsonProperty]
public DateTimeOffset PickupTime { get; set; }
[JsonIgnore]
public string PartitionKey
{
get => $"{Medallion}_{HackLicense}_{VendorId}";
}
Ez a tulajdonság explicit partíciókulcs megadására szolgál az Event Hubsba való küldéskor:
using (var client = pool.GetObject())
{
return client.Value.SendAsync(new EventData(Encoding.UTF8.GetBytes(
t.GetData(dataFormat))), t.PartitionKey);
}
Folyamfeldolgozás
A streamfeldolgozási feladat egy SQL-lekérdezéssel van definiálva, amely több különböző lépésből áll. Az első két lépés a két bemeneti stream rekordjait választja ki.
WITH
Step1 AS (
SELECT PartitionId,
TRY_CAST(Medallion AS nvarchar(max)) AS Medallion,
TRY_CAST(HackLicense AS nvarchar(max)) AS HackLicense,
VendorId,
TRY_CAST(PickupTime AS datetime) AS PickupTime,
TripDistanceInMiles
FROM [TaxiRide] PARTITION BY PartitionId
),
Step2 AS (
SELECT PartitionId,
medallion AS Medallion,
hack_license AS HackLicense,
vendor_id AS VendorId,
TRY_CAST(pickup_datetime AS datetime) AS PickupTime,
tip_amount AS TipAmount
FROM [TaxiFare] PARTITION BY PartitionId
),
A következő lépés összekapcsolja a két bemeneti streamet az egyes streamek megfelelő rekordjainak kiválasztásához.
Step3 AS (
SELECT tr.TripDistanceInMiles,
tf.TipAmount
FROM [Step1] tr
PARTITION BY PartitionId
JOIN [Step2] tf PARTITION BY PartitionId
ON tr.PartitionId = tf.PartitionId
AND tr.PickupTime = tf.PickupTime
AND DATEDIFF(minute, tr, tf) BETWEEN 0 AND 15
)
Ez a lekérdezés olyan mezők rekordjait illeszti össze, amelyek egyedileg azonosítják az egyező rekordokat (PartitionId és PickupTime).
Feljegyzés
Azt szeretnénk, hogy a TaxiRide és TaxiFare streamek találkozzanak az egyedi kombinációval, amely tartalmazza a Medallion, HackLicense, VendorId és PickupTime elemeket. Ebben az esetben a PartitionId mezőket és Medallion a HackLicensemezőket VendorId fedi le, de általában nem szabad ezt figyelembe venni.
A Stream Analyticsben az illesztések időbeliek, ami azt jelenti, hogy a rekordok egy adott időkereten belül lesznek összekapcsolva. Ellenkező esetben előfordulhat, hogy a feladatnak határozatlan ideig várnia kell egyezésre. A DATEDIFF függvény azt határozza meg, hogy egy egyezéshez mennyi ideig lehet két egyező rekordot elválasztani.
A feladat utolsó lépése az átlagos borravalót mérföldenként számítja ki, ötperces időablak szerint csoportosítva.
SELECT System.Timestamp AS WindowTime,
SUM(tr.TipAmount) / SUM(tr.TripDistanceInMiles) AS AverageTipPerMile
INTO [TaxiDrain]
FROM [Step3] tr
GROUP BY HoppingWindow(Duration(minute, 5), Hop(minute, 1))
A Stream Analytics számos ablakfüggvényt biztosít. A felugró ablak előrehalad egy meghatározott időszakkal, ebben az esetben ugrásonként egy perccel. Az eredmény az elmúlt öt perc mozgóátlagának kiszámítása.
Az itt látható architektúra csak a Stream Analytics-feladat eredményeit menti az Azure Cosmos DB-be. Big Data-forgatókönyv esetén érdemes lehet az Event Hubs Capture használatával is menteni a nyers eseményadatokat az Azure Blob Storage-ba. A nyers adatok megőrzése lehetővé teszi, hogy kötegelt lekérdezéseket futtasson az előzményadatokon később, hogy új megállapításokat nyerhessen az adatokból.
Megfontolások
Ezek a szempontok implementálják az Azure Well-Architected Framework alappilléreit, amely a számítási feladatok minőségének javítására használható vezérelvek halmaza. További információ: Microsoft Azure Well-Architected Framework.
Költségoptimalizálás
A költségoptimalizálás a szükségtelen kiadások csökkentésének és a működési hatékonyság javításának módjairól szól. További információ: Költségoptimalizálásitervezési felülvizsgálati ellenőrzőlistája.
Az Azure díjkalkulátorával megbecsülheti költségeit. Íme néhány szempont a referenciaarchitektúrában használt szolgáltatásokra vonatkozóan.
Azure Stream Analytics
Az Azure Stream Analytics ára az adatok szolgáltatásba történő feldolgozásához szükséges streamelési egységek száma (0,11 USD/óra).
A Stream Analytics költséges lehet, ha nem valós időben dolgozza fel az adatokat, vagy csak kis mennyiségben végzi az adatok feldolgozását. Ezekben a használati esetekben érdemes lehet az Azure Functions vagy a Logic Apps használatával adatokat áthelyezni az Azure Event Hubsból egy adattárba.
Az Azure Event Hubs és az Azure Cosmos DB
Az Azure Event Hubs és az Azure Cosmos DB költségeivel kapcsolatos szempontokért tekintse meg az Azure Databricks referenciaarchitektúrájával végzett Stream-feldolgozást .
Működési kiválóság
Az Operational Excellence az üzemeltetési folyamatok összességét jelenti, amelyek az alkalmazások üzembe helyezését és azok folyamatos futtatását biztosítják éles környezetben. További információ: Működési kiválóságitervezési felülvizsgálati ellenőrzőlistája.
Figyelés
Bármilyen streamfeldolgozási megoldás esetén fontos a rendszer teljesítményének és állapotának monitorozása. Az Azure Monitor metrikákat és diagnosztikai naplókat gyűjt az architektúrában használt Azure-szolgáltatásokhoz. Az Azure Monitor az Azure platformba van beépítve, és nem igényel további kódot az alkalmazásban.
Az alábbi figyelmeztető jelek bármelyike azt jelzi, hogy fel kell méreteznie a megfelelő Azure-erőforrást:
- Az Event Hubs szabályozza a kérelmeket, vagy közel van a napi üzenetkvótához.
- A Stream Analytics-feladat következetesen a lefoglalt streamegységek (SU) több mint 80%-át használja.
- Az Azure Cosmos DB megkezdi a kérések szabályozását.
A referenciaarchitektúra tartalmaz egy egyéni irányítópultot, amely az Azure Portalon van üzembe helyezve. Az architektúra üzembe helyezése után megtekintheti az irányítópultot az Azure Portal megnyitásával és az irányítópultok listájának kiválasztásávalTaxiRidesDashboard. Az egyéni irányítópultok Azure Portalon való létrehozásáról és üzembe helyezéséről további információt az Azure-irányítópultok programozott létrehozása című témakörben talál.
Az alábbi képen az irányítópult látható, miután a Stream Analytics-feladat körülbelül egy órán át futott.
A bal alsó panelen látható, hogy a Stream Analytics-feladat SU-felhasználása az első 15 percben növekszik, majd stabilizálódik. Ez egy tipikus minta, mivel a feladat stabil állapotba kerül.
Vegye észre, hogy az Event Hubs szabályozza a kérelmeket, amint az a jobb felső panelen látható. Az időnként szabályozott kérések nem okoznak problémát, mert az Event Hubs kliens SDK automatikusan újrapróbálkozik, amikor szabályozási hibát kap. Ha azonban konzisztens szabályozási hibákat lát, az azt jelenti, hogy az eseményközpontnak több átviteli egységre van szüksége. Az alábbi grafikon egy tesztfuttatást mutat be az Event Hubs automatikus felfújási funkciójával, amely szükség szerint automatikusan skálázza ki az átviteli egységeket.
Az automatikus felfújás körülbelül 06:35-kor lett engedélyezve. A korlátozott kérelmeknél látható a "p" érték csökkenése, mivel az Event Hubs automatikusan fel lett skálázva 3 adategységre.
Érdekes módon ez növelte a SU-kihasználtságot a Stream Analytics feladatban. A szabályozással az Event Hubs mesterségesen csökkentette a Stream Analytics-feladat adatbeviteli sebességét. Valójában gyakori, hogy az egyik teljesítménybeli szűk keresztmetszet feloldása egy másikat fed fel. Ebben az esetben a Stream Analytics feladathoz szükséges további SU kiosztása megoldotta a problémát.
DevOps
Hozzon létre külön erőforráscsoportokat éles, fejlesztési és tesztelési környezetekhez. A külön erőforráscsoportok használata megkönnyíti az üzemelő példányok felügyeletét, a tesztkörnyezetek törlését és a hozzáférési jogok kiosztását.
Az Azure Resource Manager sablon segítségével telepítheti az Azure erőforrásokat, az infrastruktúra mint kód (IaC) folyamatát követve. Sablonokkal egyszerűbb automatizálni az üzembe helyezéseket az Azure DevOps Services vagy más CI/CD-megoldások használatával.
Helyezze az egyes számítási feladatokat egy külön üzembehelyezési sablonba, és tárolja az erőforrásokat a forrásvezérlő rendszerekben. A sablonokat egy CI/CD-folyamat részeként együtt vagy egyenként is üzembe helyezheti, így egyszerűbbé teheti az automatizálási folyamatot.
Ebben az architektúrában az Azure Event Hubs, a Log Analytics és az Azure Cosmos DB egyetlen számítási feladatként van azonosítva. Ezeket az erőforrásokat egyetlen ARM-sablon tartalmazza.
Fontolja meg a számítási feladatok fokozatos előkészítését. Helyezze üzembe különböző szakaszokban, és futtasson validációs ellenőrzéseket minden szakaszban, mielőtt továbblép a következőre. Így nagyon ellenőrzött módon telepítheti a frissítéseket az éles környezetekbe, és minimalizálhatja a nem várt üzembe helyezési problémákat.
Fontolja meg az Azure Monitor használatát a streamfeldolgozási folyamat teljesítményének elemzéséhez.
További információ: a Microsoft Azure Well-Architected Framework működési kiválósági pillére.
Teljesítményhatékonyság
A teljesítményhatékonyság az a képesség, hogy a számítási feladatok skálázhatók, hogy hatékonyan megfeleljenek a felhasználók által támasztott követelményeknek. További információ: Teljesítményhatékonyságtervezési felülvizsgálati ellenőrzőlistája.
Event Hubs
Az Event Hubs átviteli kapacitását átviteli egységekben mérik. Az eseményközpontok automatikus méretezéséhez engedélyezze az automatikus felfújást, amely automatikusan skálázza az átviteli egységeket a forgalom alapján, egy konfigurált maximális értékre.
Stream Analytics
A Stream Analytics esetében a feladathoz lefoglalt számítási erőforrásokat streamelési egységekben mérik. A Stream Analytics-feladatok akkor méretezhetők a legjobban, ha a feladat párhuzamosítható. Így a Stream Analytics több számítási csomóponton is elosztja a feladatot.
Az Event Hubs-bemenethez használja a PARTITION BY kulcsszót a Stream Analytics-feladat particionálásához. Az adatok részhalmazokra lesznek osztva az Event Hubs-partíciók alapján.
Az ablakfüggvények és az időbeli illesztések további su-t igényelnek. Ha lehetséges, használja PARTITION BY az egyes partíciók külön-külön történő feldolgozását. További információ: Streamelési egységek értelmezése és módosítása.
Ha nem lehetséges a teljes Stream Analytics-feladat párhuzamosítása, próbálja meg több lépésre bontani a feladatot, kezdve egy vagy több párhuzamos lépéssel. Így az első lépések párhuzamosan is futtathatók. Ebben a referenciaarchitektúrában például:
- Az 1. és a 2
SELECT. lépés olyan utasítások, amelyek egyetlen partíción belül választják ki a rekordokat. - A 3. lépés particionált illesztéseket hajt végre két bemeneti adatfolyamon. Ez a lépés kihasználja azt a tényt, hogy az egyező rekordok ugyanazt a partíciókulcsot használják, és így garantáltan ugyanazzal a partícióazonosítóval rendelkeznek az egyes bemeneti adatfolyamokban.
- A 4. lépés összesíti az összes partíciót. Ez a lépés nem párhuzamosítható.
A Stream Analytics-feladatdiagram használatával megtekintheti, hogy hány partíció van hozzárendelve a feladat egyes lépéseihez. Az alábbi diagram a referenciaarchitektúra feladatábraét mutatja be:
Azure Cosmos DB
Az Azure Cosmos DB átviteli kapacitásának mérése kérelemegységekben (RU) történik. Minden Azure Cosmos DB tárolóhoz szükség van egy partition kulcsra, és minden dokumentumnak tartalmaznia kell ezt a kulcsot. Egyetlen fizikai partíció legfeljebb 10 000 RU/s kezelésére képes, ezért az Azure Cosmos DB a partíciókulcs alapján több fizikai partíció között osztja el az adatokat, hogy ezen a korláton túl is lehessen skálázni. Válasszon olyan partíciókulcsot, amely egyenletesen osztja el a tárolást és a kérelmek mennyiségét, hogy elkerülje a forró partíciókat.
Ebben a referenciaarchitektúrában az új dokumentumok percenként csak egyszer jönnek létre (az ugróablak-intervallumon belül), így az áteresztőképességi követelmények alacsonyak. Ennek ellenére válasszon ki egy partíciókulcsot, amely támogatja a várt lekérdezési mintákat és növekedést.