Parsování dat JSON a Avro v Azure Stream Analytics

Služba Azure Stream Analytics podporuje zpracování událostí ve formátech dat CSV, JSON a Avro. Data JSON i Avro můžou být strukturovaná a obsahují některé komplexní typy, jako jsou vnořené objekty (záznamy) a pole.

Záznam datových typů

Datové typy záznamů se používají k reprezentaci polí JSON a Avro, pokud se ve vstupních datových proudech používají odpovídající formáty. Tyto příklady ukazují ukázkový senzor, který čte vstupní události ve formátu JSON. Tady je příklad jedné události:

{
    "DeviceId" : "12345",
    "Location" :
    {
        "Lat": 47,
        "Long": 122
    },
    "SensorReadings" :
    {
        "Temperature" : 80,
        "Humidity" : 70,
        "CustomSensor01" : 5,
        "CustomSensor02" : 99,
        "SensorMetadata" : 
        {
        "Manufacturer":"ABC",
        "Version":"1.2.45"
        }
    }
}

Přístup k vnořeným polím ve známém schématu

Pro přístup k vnořeným polím přímo z dotazu použijte zápis tečky (.). Tento dotaz například vybere souřadnice zeměpisné šířky a délky ve vlastnosti Umístění v předchozích datech JSON. Pomocí zápisu tečky můžete procházet více úrovní, jak je znázorněno v následujícím fragmentu kódu:

SELECT
    DeviceID,
    Location.Lat,
    Location.Long,
    SensorReadings.Temperature,
    SensorReadings.SensorMetadata.Version
FROM input

Výsledkem je:

|DeviceID|Lat|Long|Temperature|Version|
|-|-|-|-|-|
|12345|47|122|80|1.2.45|

Vybrat všechny vlastnosti

Pomocí zástupného znaku můžete vybrat všechny vlastnosti vnořeného záznamu * . Představte si následující příklad:

SELECT
    DeviceID,
    Location.*
FROM input

Výsledkem je:

|DeviceID|Lat|Long|
|-|-|-|
|12345|47|122|

Přístup k vnořeným polím, pokud název vlastnosti je proměnnou

Pokud je název vlastnosti proměnnou, použijte funkci GetRecordPropertyValue . Tato funkce vám pomůže sestavovat dynamické dotazy bez pevně zakódovaných názvů vlastností.

Představte si například, že se ukázkový datový proud musí spojit s referenčními daty obsahujícími prahové hodnoty pro každý senzor zařízení. Fragment těchto referenčních dat se zobrazí v následujícím fragmentu kódu.

{
    "DeviceId" : "12345",
    "SensorName" : "Temperature",
    "Value" : 85
},
{
    "DeviceId" : "12345",
    "SensorName" : "Humidity",
    "Value" : 65
}

Cílem je sloučit ukázkovou datovou sadu z horní části článku s referenčními daty a generovat jednu událost pro každé měření senzoru přesahující jeho prahovou hodnotu. Toto spojení znamená, že jedna událost může generovat více výstupních událostí, pokud více senzorů překročí příslušné prahové hodnoty. Pokud chcete dosáhnout podobných výsledků bez spojení, podívejte se na následující příklad:

SELECT
    input.DeviceID,
    thresholds.SensorName,
    "Alert: Sensor above threshold" AS AlertMessage
FROM input      -- stream input
JOIN thresholds -- reference data input
ON
    input.DeviceId = thresholds.DeviceId
WHERE
    GetRecordPropertyValue(input.SensorReadings, thresholds.SensorName) > thresholds.Value

GetRecordPropertyValue vybere vlastnost v SensorReadings , která odpovídá názvu vlastnosti pocházející z referenčních dat. Pak extrahuje přidruženou hodnotu ze SensorReadings.

Výsledkem je:

|DeviceID|SensorName|AlertMessage|
| - | - | - |
| 12345 | Humidity | Alert: Sensor above threshold |

Převod polí záznamů na samostatné události

Chcete-li převést pole záznamů na samostatné události, použijte operátor APPLY společně s funkcí GetRecordProperties .

Pomocí původních ukázkových dat můžete pomocí následujícího dotazu extrahovat vlastnosti do různých událostí:

SELECT
    event.DeviceID,
    sensorReading.PropertyName,
    sensorReading.PropertyValue
FROM input as event
CROSS APPLY GetRecordProperties(event.SensorReadings) AS sensorReading

Výsledkem je:

|DeviceID|SensorName|AlertMessage|
|-|-|-|
|12345|Temperature|80|
|12345|Humidity|70|
|12345|CustomSensor01|5|
|12345|CustomSensor02|99|
|12345|SensorMetadata|[object Object]|

Pomocí funkce WITH můžete tyto události směrovat do různých cílů:

WITH Stage0 AS
(
    SELECT
        event.DeviceID,
        sensorReading.PropertyName,
        sensorReading.PropertyValue
    FROM input as event
    CROSS APPLY GetRecordProperties(event.SensorReadings) AS sensorReading
)

SELECT DeviceID, PropertyValue AS Temperature INTO TemperatureOutput FROM Stage0 WHERE PropertyName = 'Temperature'
SELECT DeviceID, PropertyValue AS Humidity INTO HumidityOutput FROM Stage0 WHERE PropertyName = 'Humidity'

Parsování záznamu JSON v referenčních datech SQL

Když ve své úloze použijete Azure SQL Database jako referenční data, můžete zahrnout sloupec, který obsahuje data ve formátu JSON. Následující příklad ukazuje tento formát:

|DeviceID|Data|
|-|-|
|12345|{"key": "value1"}|
|54321|{"key": "value2"}|

Záznam JSON ve sloupci Data můžete analyzovat napsáním jednoduché uživatelem definované funkce JavaScriptu.

function parseJson(string) {
return JSON.parse(string);
}

Pokud chcete získat přístup k polím záznamů JSON, vytvořte krok v dotazu Stream Analytics, jak je znázorněno v následujícím příkladu.

WITH parseJson as
(
SELECT DeviceID, udf.parseJson(sqlRefInput.Data) as metadata,
FROM sqlRefInput
)

SELECT metadata.key
INTO output
FROM streamInput
JOIN parseJson 
ON streamInput.DeviceID = parseJson.DeviceID

Pole datových typů

Datové typy pole jsou seřazenou kolekcí hodnot. Tato část podrobně popisuje některé typické operace s hodnotami pole. Tyto příklady používají funkce GetArrayElement, GetArrayElements, GetArrayLength a APPLY operátor.

Tady je příklad události. Oba CustomSensor03 a SensorMetadata jsou typu pole:

{
    "DeviceId" : "12345",
    "SensorReadings" :
    {
        "Temperature" : 80,
        "Humidity" : 70,
        "CustomSensor01" : 5,
        "CustomSensor02" : 99,
        "CustomSensor03": [12,-5,0]
     },
    "SensorMetadata":[
        {          
            "smKey":"Manufacturer",
            "smValue":"ABC"                
        },
        {
            "smKey":"Version",
            "smValue":"1.2.45"
        }
    ]
}

Práce s konkrétním prvkem pole

Vyberte prvek pole v zadaném indexu (vyberte první prvek pole):

SELECT
    GetArrayElement(SensorReadings.CustomSensor03, 0) AS firstElement
FROM input

Výsledkem je:

|firstElement|
|-|
|12|

Vybrat délku pole

SELECT
    GetArrayLength(SensorReadings.CustomSensor03) AS arrayLength
FROM input

Výsledkem je:

|arrayLength|
|-|
|3|

Převod prvků pole na samostatné události

Vyberte všechny prvky pole jako jednotlivé události. Operátor APPLY společně s integrovanou funkcí GetArrayElements extrahuje všechny prvky pole jako jednotlivé události:

SELECT
    DeviceId,
	CustomSensor03Record.ArrayIndex,
	CustomSensor03Record.ArrayValue
FROM input
CROSS APPLY GetArrayElements(SensorReadings.CustomSensor03) AS CustomSensor03Record

Výsledkem je:

|DeviceId|ArrayIndex|ArrayValue|
|-|-|-|
|12345|0|12|
|12345|1|-5|
|12345|2|0|
SELECT   
    i.DeviceId,	
    SensorMetadataRecords.ArrayValue.smKey as smKey,
    SensorMetadataRecords.ArrayValue.smValue as smValue
FROM input i
CROSS APPLY GetArrayElements(SensorMetadata) AS SensorMetadataRecords

Výsledkem je:

|DeviceId|smKey|smValue|
|-|-|-|
|12345|Manufacturer|ABC|
|12345|Version|1.2.45|

Pokud chcete zobrazit extrahovaná pole ve sloupcích, překlopte datovou sadu pomocí syntaxe WITH spolu s operací JOIN . Toto spojení vyžaduje podmínku časové hranice , která brání duplikaci:

WITH DynamicCTE AS (
	SELECT   
		i.DeviceId,
		SensorMetadataRecords.ArrayValue.smKey as smKey,
		SensorMetadataRecords.ArrayValue.smValue as smValue
	FROM input i
	CROSS APPLY GetArrayElements(SensorMetadata) AS SensorMetadataRecords 
)

SELECT
	i.DeviceId,
	i.Location.*,
	V.smValue AS 'smVersion',
	M.smValue AS 'smManufacturer'
FROM input i
LEFT JOIN DynamicCTE V ON V.smKey = 'Version' and V.DeviceId = i.DeviceId AND DATEDIFF(minute,i,V) BETWEEN 0 AND 0 
LEFT JOIN DynamicCTE M ON M.smKey = 'Manufacturer' and M.DeviceId = i.DeviceId AND DATEDIFF(minute,i,M) BETWEEN 0 AND 0

Výsledkem je:

|DeviceId|Lat|Long|smVersion|smManufacturer|
|-|-|-|-|-|
|12345|47|122|1.2.45|ABC|

Data Types in Azure Stream Analytics