Creare un flusso di eventi con una destinazione di inserimento diretto della eventhouse usando le API

Questo articolo fornisce una guida dettagliata per creare un flusso di eventi con una destinazione DirectIngestion eventhouse usando le API.

Completare tre passaggi:

Prerequisiti

  • Hai accesso a un'area di lavoro con capacità Fabric o di tipo Versione di valutazione di Fabric, con autorizzazioni di Collaboratore o di livello superiore.

Requisiti di autenticazione e autorizzazione

Per usare Fabric API, ottenere prima un token di Microsoft Entra ID per il servizio Fabric e quindi usare tale token nell'intestazione Authorization della chiamata API. È possibile ottenere un token di Fabric usando MSAL.NET.

Ottenere il token usando MSAL.NET

Se l'applicazione deve accedere alle API di Fabric usando un'entità servizio, è possibile usare la libreria MSAL.NET per acquisire un token di accesso. Seguire la guida introduttiva all'API Fabric per creare un'app console C#, che acquisisce un token Microsoft Entra ID usando la libreria MSAL.NET e quindi usa C# HttpClient per chiamare l'API List Workspaces.

Note

Se il flusso di eventi creato include tutte le origini che usano una connessione cloud, assicurarsi che l'identità usata per ottenere il token disponga dell'autorizzazione per accedere a tale connessione cloud, sia che si tratti di un'entità servizio o di un utente.

Passaggio 1: Creare eventhouse per API

Indirizzo API e parametri

POST https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/eventhouses
Parametro In Obbligatorio Descrzione
workspaceId Percorso Yes Area di lavoro in cui viene creato l'elemento Eventhouse.

Token

Authorization: Bearer <fabric_access_token>
Content-Type: application/json

Payload

Usare un payload minimo per creare l'elemento Eventhouse.

{
  "displayName": "es-eh-demo"
}

Se è necessario effettuare il provisioning avanzato in base alle parti di definizione, usare il contratto di definizione Eventhouse nella documentazione di riferimento.

Esempio di risposta

{
  "id": "00000000-0000-0000-0000-000000000000",
  "type": "Eventhouse",
  "displayName": "es-eh-demo",
  "description": "",
  "workspaceId": "00000000-0000-0000-0000-000000000000"
}

Acquisisci il valore per le fasi successive:

  • ID dell'elemento Eventhouse (id)

Passaggio 2: Creare un database KQL con tabella e mappatura nella definizione

Usare l'API REST Fabric per creare il database KQL con la relativa definizione di schema in una singola chiamata. Questo passaggio usa solo gli endpoint API Fabric (non sono necessarie chiamate agli endpoint di gestione di Eventhouse). Il nuovo database viene aggiunto automaticamente a Eventhouse a partire dal passaggio 1 tramite parentEventhouseItemId, e la tabella e il mapping di acquisizione sono definiti in DatabaseSchema.kql.

Indirizzo API e parametri

POST https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/kqlDatabases
Parametro In Obbligatorio Descrzione
workspaceId Percorso Yes ID dell'area di lavoro in cui è stata creata la eventhouse nel passaggio 1. Deve essere la stessa area di lavoro.

Token

Authorization: Bearer <fabric_access_token>
Content-Type: application/json

Parti della definizione

Preparare due parti di definizione:

  • DatabaseProperties.json: collega il database all'Eventhouse dal passaggio 1 utilizzando il campo parentEventhouseItemId (questa è la connessione chiave tra il database e l'Eventhouse).
  • DatabaseSchema.kql: definisce e crea la struttura della tabella KQL e il mapping di inserimento che verranno eseguiti automaticamente al momento della creazione del database.

DatabaseProperties.json esempio:

{
  "databaseType": "ReadWrite",
  "parentEventhouseItemId": "<eventhouseItemId>",
  "oneLakeCachingPeriod": "P36500D",
  "oneLakeStandardStoragePeriod": "P36500D"
}

DatabaseSchema.kql esempio:

.create-merge table Orders (id:string, eventTime:datetime, amount:real)
.create-or-alter table Orders ingestion json mapping 'orders_json_map' "[{\"column\":\"id\",\"Properties\":{\"path\":\"$.id\"}},{\"column\":\"eventTime\",\"Properties\":{\"path\":\"$.eventTime\"}},{\"column\":\"amount\",\"Properties\":{\"path\":\"$.amount\"}}]"

Codificare parti di definizione in Base64

$databaseProperties = @'
{
  "databaseType": "ReadWrite",
  "parentEventhouseItemId": "<eventhouseItemId>",
  "oneLakeCachingPeriod": "P36500D",
  "oneLakeStandardStoragePeriod": "P36500D"
}
'@

$databaseSchema = @'
.create-merge table Orders (id:string, eventTime:datetime, amount:real)
.create-or-alter table Orders ingestion json mapping 'orders_json_map' "[{\"column\":\"id\",\"Properties\":{\"path\":\"$.id\"}},{\"column\":\"eventTime\",\"Properties\":{\"path\":\"$.eventTime\"}},{\"column\":\"amount\",\"Properties\":{\"path\":\"$.amount\"}}]"
'@

$base64DatabaseProperties = [Convert]::ToBase64String([System.Text.Encoding]::UTF8.GetBytes($databaseProperties))
$base64DatabaseSchema = [Convert]::ToBase64String([System.Text.Encoding]::UTF8.GetBytes($databaseSchema))

Payload della richiesta

{
  "displayName": "es-kql-demo",
  "description": "KQL database created by API with schema definition",
  "definition": {
    "parts": [
      {
        "path": "DatabaseProperties.json",
        "payload": "<base64DatabaseProperties>",
        "payloadType": "InlineBase64"
      },
      {
        "path": "DatabaseSchema.kql",
        "payload": "<base64DatabaseSchema>",
        "payloadType": "InlineBase64"
      }
    ]
  }
}

Risposta

L'API restituisce 202 Accepted con un'intestazione Location per questa operazione a esecuzione prolungata. Eseguire il polling dell'endpoint nell'intestazione Location per verificare il completamento dell'operazione.

Quando l'operazione ha esito positivo, il database KQL viene creato con la tabella Orders e il mapping di inserimento orders_json_map JSON già configurato.

Acquisire questi valori per il passaggio 3:

  • tableName: Orders
  • mappingRuleName: orders_json_map

Passaggio 3: Creare eventstream in modalità DirectIngestion

Indirizzo API e parametri

POST https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items
Parametro In Obbligatorio Descrzione
workspaceId Percorso Yes Area di lavoro in cui viene creato l'elemento Eventstream.

Payload della topologia Eventstream

Questo payload di esempio usa SampleData come origine e Eventhouse DirectIngestion come destinazione. Mantenere le proprietà di destinazione allineate alla versione dell'API.

La topologia definisce tre componenti:

Componente Purpose
sources Fonte dei dati di input (in questo esempio, un feed di esempio di dati del mercato azionario)
streams Pipeline che instrada i dati dalle origini alle destinazioni, inclusi i flussi predefiniti e derivati
destinations Destinazione di output in cui confluiscono i dati (in questo caso, l'eventhouse in modalità di inserimento diretto)

Assicurati che itemId e workspaceId nella destinazione corrispondano all'Eventhouse del passaggio 1 e che tableName e mappingRuleName corrispondano a ciò che hai creato nel passaggio 2.

{
  "sources": [
    {
      "name": "sample-data-source",
      "type": "SampleData",
      "properties": {
        "type": "StockMarket"
      }
    }
  ],
  "destinations": [
    {
      "name": "eventhouse-direct-ingestion",
      "type": "Eventhouse",
      "properties": {
        "dataIngestionMode": "DirectIngestion",
        "workspaceId": "<eventhouseWorkspaceId>",
        "itemId": "<eventhouseItemId>",
        "tableName": "Orders",
        "connectionName": "es-eh-conn-7f3a",
        "mappingRuleName": "orders_json_map"
      },
      "inputNodes": [
        {
          "name": "eventstream-main-stream"
        }
      ]
    }
  ],
  "streams": [
    {
      "name": "eventstream-main-stream",
      "type": "DefaultStream",
      "properties": {},
      "inputNodes": [
        {
          "name": "sample-data-source"
        }
      ]
    }
  ],
  "operators": [],
  "compatibilityLevel": "1.1"
}

Campi di destinazione usati in modalità DirectIngestion:

Campo Origine del valore
workspaceId ID dell'area di lavoro in cui è stata creata la eventhouse nel passaggio 1
itemId L'ID dell'elemento Eventhouse (id) restituito nella risposta del passaggio 1
connectionName Qualsiasi nome univoco fino a 40 caratteri. È consigliabile usare un suffisso casuale, ad esempio es-eh-conn-7f3a.
tableName Nome della tabella del passaggio 2
mappingRuleName Nome della regola di mappatura dal Passaggio 2

Codificare la topologia in Base64

$json = Get-Content -Path "eventstream.json" -Raw
$bytes = [System.Text.Encoding]::UTF8.GetBytes($json)
$base64Eventstream = [Convert]::ToBase64String($bytes)

Esempio di richiesta

La richiesta di esempio include un payload con due parti di definizione con codifica Base64: eventstream.json (la topologia definita in precedenza) e .platform (file di metadati, necessario per tutti gli elementi Fabric).

{
  "displayName": "es-directingest-demo",
  "type": "Eventstream",
  "description": "Eventstream created by API in DirectIngestion mode",
  "definition": {
    "parts": [
      {
        "path": "eventstream.json",
        "payload": "<base64Eventstream>",
        "payloadType": "InlineBase64"
      },
      {
        "path": ".platform",
        "payload": "<base64Platform>",
        "payloadType": "InlineBase64"
      }
    ]
  }
}

Per generare un valore di payload personalizzato .platform , eseguire le operazioni seguenti:

  1. Creare un .platform file usando il formato illustrato in Creare un elemento Eventstream con definizione.
  2. Codifica in base64 l'intero .platform contenuto del file usando lo stesso approccio illustrato in Codificare la topologia in Base64.
  3. Usare la stringa codificata come valore per il payload campo, sostituendo <base64Platform> nella richiesta di esempio.

Example:

{
  "path": ".platform",
  "payload": "ewogICIkc2NoZW1hIjogImh0dHBzOi8vZGV2ZWxvcGVyLm1pY3Jvc29mdC5jb20vanNvbi1zY2hlbWFzL2ZhYnJpYy9naXRJbnRlZ3JhdGlvbi9wbGF0Zm9ybVByb3BlcnRpZXMvMi4wLjAvc2NoZW1hLmpzb24iLAogICJtZXRhZGF0YSI6IHsKICAgICJ0eXBlIjogIkV2ZW50c3RyZWFtIiwKICAgICJkaXNwbGF5TmFtZSI6ICJhbGV4LWVzMSIKICB9LAogICJjb25maWciOiB7CiAgICAidmVyc2lvbiI6ICIyLjAiLAogICAgImxvZ2ljYWxJZCI6ICIwMDAwMDAwMC0wMDAwLTAwMDAtMDAwMC0wMDAwMDAwMDAwMDAiCiAgfQp9",
  "payloadType": "InlineBase64"
}

Risposta di esempio

202 Accepted

L'API restituisce 202 Accepted per questa operazione a esecuzione prolungata. A seconda del client, il corpo della risposta potrebbe essere vuoto o potrebbe contenere un valore letterale null. Se la risposta include un'intestazione Location, usarla per verificare tramite polling il completamento dell'operazione.

Elenco di controllo end-to-end

  1. Ottenere un token di Fabric (aud = https://api.fabric.microsoft.com).
  2. Crea Eventhouse e acquisisci workspaceId e itemId.
  3. Creare un database KQL con DatabaseProperties.json e DatabaseSchema.kql.
  4. Compila e codifica in Base64 eventstream.json.
  5. Creare un Eventstream con una destinazione DirectIngestion di Eventhouse.

References