Samouczek: tworzenie mapy w czasie rzeczywistym przy użyciu interfejsów API REST i Python

Fabric Maps może wizualizować dane geoprzestrzenne w czasie rzeczywistym, łącząc się z zestawami danych eventhouse, które są stale aktualizowane poprzez pozyskiwanie danych z usługi Eventstream.

W przeciwieństwie do scenariuszy statycznych korzystających z plików przechowywanych w usłudze Lakehouse, w tym samouczku przedstawiono architekturę opartą na zdarzeniach przesyłania strumieniowego, w której:

  • Zdarzenia są przyjmowane do hurtowni zdarzeń
  • Dane są odpytywane przy użyciu języka Kusto Query Language (KQL)
  • Mapa dynamicznie odświeża się po nadejściu nowych danych

Ten samouczek koncentruje się na automatyzacji kompleksowego przepływu pracy od początku do końca przy użyciu interfejsów API REST usługi Fabric i języka Python, dzięki czemu można programowo aprowizować zasoby i konfigurować obsługę mapy w czasie rzeczywistym. W przypadku scenariuszy danych statycznych korzystających z plików lakehouse zobacz Utwórz mapę statyczną przy użyciu interfejsów API REST i Python.

Z tego samouczka dowiesz się, jak tworzyć i automatyzować rozwiązanie geoprzestrzenne w czasie rzeczywistym w Microsoft Fabric przy użyciu usługi Eventstream, Eventhouse i KQL.

Korzystając z interfejsu API REST usługi Fabric:

  • Utwórz eventhouse i bazę danych KQL
  • Utwórz eventstream, aby importować dane do eventhouse
  • Utwórz mapę z definicją w tekście, która odwołuje się do danych Eventhouse
  • Skonfiguruj warstwę mapową z okresowym odświeżaniem, aby uzyskiwać aktualizacje w czasie rzeczywistym
  • Zainicjuj początkowe zdarzenia, aby mapa od razu wyświetlała dane.

Aby symulować ciągłe przesyłanie strumieniowe i obserwować aktualizację mapy niemal w czasie rzeczywistym, najpierw ukończ ten samouczek, a następnie przejdź do kolejnego samouczka Samouczek: symulowanie pozyskiwania danych w czasie rzeczywistym dla mapy przy użyciu interfejsów API REST i języka Python, który opiera się bezpośrednio na magazynie zdarzeń, strumieniu zdarzeń, funkcji KQL i mapie utworzonych tutaj.

Omówienie scenariusza: śledzenie zasobów w czasie rzeczywistym

Ten samouczek opiera się na scenariuszu śledzenia zasobów w czasie rzeczywistym, podobnym do scenariusza śledzenia floty użytego w oryginalnym samouczku Fabric Maps Samouczek: Tworzenie trasowania zleceń roboczych w czasie rzeczywistym za pomocą usługi Fabric Maps.

W tym scenariuszu:

  • Pojazdy okresowo emitują aktualizacje lokalizacji
  • Zdarzenia lokalizacyjne są importowane do magazynu zdarzeń
  • Mapa wyświetla najnowsze pozycje pojazdów i aktualizacje automatycznie po nadejściu nowych zdarzeń

Ten wzorzec jest reprezentatywny dla typowych przypadków użycia operacji w czasie rzeczywistym, takich jak:

  • Śledzenie floty
  • Wysyłka zleceń roboczych
  • Monitorowanie zasobów i sprzętu

Microsoft Fabric używa Eventstream i Eventhouse pozyskiwania, przetwarzania i analizowania danych przesyłanych strumieniowo niemal w czasie rzeczywistym, dzięki czemu można wizualizować dane operacyjne na żywo bezpośrednio na mapie.

Ten samouczek przedstawia typowy wzorzec automatyzacji w usłudze Fabric: tworzenie infrastruktury → pozyskiwanie danych strumieniowych → weryfikacja poprawności pozyskiwania → renderowanie mapy.

Prerequisites

  • Python 3.10 lub nowszy
  • Azure CLI
  • Identyfikator obszaru roboczego Fabric
  • Uprawnienia do wywoływania interfejsów API REST Fabric, takich jak:
    • Item.ReadWrite.All

Note

Delegowane zakresy, takie jak Item.ReadWrite.All, są przyznawane zalogowanej tożsamości poprzez jej rolę obszaru roboczego. Upewnij się, że tożsamość używana z az login ma przypisaną rolę Contributor, Member lub Admin w docelowym obszarze roboczym Fabric przed uruchomieniem skryptu.

Uwierzytelnianie

W tym samouczku jest używany program DefaultAzureCredential, który może uwierzytelniać się przy użyciu kilku źródeł poświadczeń lokalnych/deweloperskich. Najprostszym podejściem dla czytelników, którzy są po raz pierwszy, jest logowanie za pomocą Azure CLI (interfejsu wiersza polecenia platformy Azure).

  1. Otwórz terminal.
  2. Run:
az login

DefaultAzureCredential może używać tożsamości zalogowanej do uzyskiwania tokenów dostępu dla:

  • Architektura sieciowa API REST (zasób: https://api.fabric.microsoft.com/.default)
  • Zapytania Kusto (KQL) w płaszczyźnie danych dla eventhouse (zasób: https://api.kusto.windows.net/.default)
  • Punkty końcowe REST usługi Power BI / Fabric używane do odpytywania operacji długotrwałych (zasób: https://analysis.windows.net/powerbi/api/.default)

Wskazówka

Informacja o https://api.fabric.microsoft.com/.default tej wartości: jest to zakres żądania tokenu, a nie adres URL, który wywołujesz bezpośrednio. Informuje firmę Microsoft Entra, że token dostępu powinien zostać wystawiony dla interfejsu API REST usługi Microsoft Fabric i powinien zawierać wszystkie uprawnienia Fabric, które zostały już przyznane uwierzytelnionej tożsamości (na przykład Item.ReadWrite.All lub Workspace.ReadWrite.All).

Zakres .default jest używany tylko podczas pozyskiwania tokenu i nigdy nie jest wysyłany do punktów końcowych interfejsu API REST Fabric.

Aby uzyskać więcej informacji na temat działania .default zakresu na platformie tożsamości firmy Microsoft, zobacz Zakresy i uprawnienia na platformie tożsamości firmy Microsoft.

Przed uruchomieniem tego samouczka zalecamy zalogowanie się do usługi Microsoft Fabric co najmniej raz:

https://app.fabric.microsoft.com

Logowanie gwarantuje, że tożsamość Fabric, członkostwo w roli obszaru roboczego i przypisania pojemności są w pełni przydzielone przed programowym uzyskaniem tokenu dostępu Microsoft Entra.

Ten krok jest szczególnie przydatny, jeśli:

  • Dopiero zaczynasz korzystać z usługi Microsoft Fabric
  • Obszar roboczy został niedawno utworzony
  • Twoje przypisanie roli zostało ostatnio dodane

Note

Ten samouczek uwierzytelnia za pomocą Entra ID firmy Microsoft poprzez DefaultAzureCredential. Interfejsy API REST sieci szkieletowej nie wymagają sesji przeglądarki, ale logowanie się do interfejsu internetowego sieci szkieletowej może zapobiec problemom z autoryzacją pierwszego uruchomienia powodowanym przez opóźnione przydzielanie ról.

Utwórz plik danych początkowych (początkowe dane mapy)

Aby mapa wyświetlała dane natychmiast po aprowizacji, skrypt wysyła do strumienia zdarzeń niewielki początkowy zestaw zdarzeń.

  1. Utwórz nowy plik w tym samym katalogu co skrypt Python: vehicle_locations_seed.csv
  2. Wklej następującą zawartość:
VehicleId,Latitude,Longitude,EventTime
V-001,47.6101,-122.3344,2026-01-01T10:00:00Z
V-002,47.6150,-122.3200,2026-01-01T10:00:00Z
V-003,47.6205,-122.3493,2026-01-01T10:00:00Z
V-004,47.6050,-122.3300,2026-01-01T10:00:00Z

Krok 1. Tworzenie nowego pliku projektu języka Python

W tym kroku utworzysz pusty plik Python, który utworzysz sekcję po sekcji.

Utwórz nowy plik o nazwie:

create_realtime_map.py

Otwórz plik w edytorze.

Krok 2. Instalowanie wymaganych bibliotek i dodawanie wymaganych instrukcji importu

W tym kroku zainstalujesz zależności i dodasz import używany przez skrypt.

Instalowanie wymaganych bibliotek

W otwartym oknie terminalu uruchom następujące polecenie:

pip install httpx azure-identity azure-eventhub

Do czego służy każda biblioteka

  • httpx: wysyła żądania HTTP do interfejsów API REST sieci szkieletowej.
  • azure-identity: zapewnia DefaultAzureCredential na potrzeby uwierzytelniania Microsoft Entra.
  • azure-eventhub: wysyła zdarzenia początkowe do punktu końcowego Eventstream zgodnego z usługą Event Hub, aby zasilić Eventhouse.

Dodawanie instrukcji importu do pliku .py

W górnej części create_realtime_map.py dodaj:

import base64
import csv
import json
import os
import time
import uuid

import httpx
from azure.eventhub import EventData, EventHubProducerClient
from azure.eventhub.exceptions import EventHubError
from azure.identity import DefaultAzureCredential

Note

EventHubError jest importowany w tym miejscu, ale nie jest używany do późniejszego użycia w skrycie. Funkcja pomocnicza seed_eventstream_from_csv obsługuje ten błąd (wraz z ConnectionError i TimeoutError) w swojej pętli ponawiania prób, dzięki czemu przejściowe błędy wysyłania do Event Hub — takie jak sytuacja, gdy niestandardowy punkt końcowy nie jest jeszcze gotowy — powodują ponowienie próby zamiast przerwania działania skryptu.

Krok 3. Dodawanie sekcji konfiguracji

W tym kroku zdefiniujesz zmienne używane przez aplikację, w tym identyfikator obszaru roboczego i nazwy zasobów.

Scentralizowanie konfiguracji w jednej Config klasie — zamiast rozpraszania wartości zakodowanych na twardo w funkcjach — daje trzy konkretne korzyści:

  • Przenośność środowiska: identyfikatory obszarów roboczych, nazwy zasobów i inne ustawienia działają w jednym miejscu, dzięki czemu można ponownie uruchomić skrypt względem innego obszaru roboczego lub maszyny, zmieniając kilka wierszy (lub zmiennej środowiskowej) zamiast wyszukiwać kod.
  • Czytelniejsze sygnatury funkcji: funkcje kroków akceptują pojedynczy obiekt cfg zamiast długich list parametrów, dzięki czemu orkiestracja w main() pozostaje czytelna.
  • Bezpieczniejsza obsługa sekretów: Poufne wartości, takie jak identyfikator przestrzeni roboczej, są wczytywane ze zmiennych środowiskowych, więc nigdy nie trafiają do repozytorium wraz ze skryptem.

Dodaj poniższe elementy pod instrukcjami import:

# =========================================================
# Configuration (centralized)
# =========================================================

class Config:
    """
    Central configuration: workspace ID, resource display names, and
    ingestion settings. A single instance is built in main() and passed
    to each step function.
    """
    def __init__(self):
        # Workspace
        self.workspace_id = os.environ.get("FABRIC_WORKSPACE_ID", "")
        if not self.workspace_id:
            raise RuntimeError("Set FABRIC_WORKSPACE_ID environment variable before running the script.")

        # Resource display names / descriptions
        self.eventhouse_display_name = "eh_realtime_locations"
        self.eventhouse_description = "Stores streaming location events for a Fabric Maps real-time tutorial"

        self.eventhouse_table_name = "VehicleLocations"
        self.kql_function_name = "LatestVehicleLocations"

        self.eventstream_display_name = "es_realtime_locations"
        self.eventstream_description = "Streams events into an eventhouse table (created via Eventstream REST API)"

        self.map_display_name = "My Real-Time Fabric Map"
        self.map_description = "Created using Fabric Maps REST API (Eventhouse + Eventstream + Kusto function)"

        # Map refresh
        self.refresh_interval_ms = 5000

        # Seed data (initial map data)
        self.seed_csv_path = os.path.join(os.path.dirname(__file__), "vehicle_locations_seed.csv")
        
        # Will be provided interactively after eventstream is created
        self.eventhub_connection_string = os.environ.get("EVENTHUB_CONNECTION_STRING", "")

Ustawianie identyfikatora obszaru roboczego przy użyciu zmiennej środowiskowej

Zamiast trwale zapisywać identyfikator obszaru roboczego bezpośrednio w skrypsie, ten samouczek odczytuje go ze zmiennej środowiskowej. Pozwala to zachować wartości specyficzne dla środowiska poza kodem źródłowym i umożliwia ponowne użycie skryptu między obszarami roboczymi lub maszynami bez jego edytowania.

Przed uruchomieniem skryptu utwórz zmienną środowiskową o nazwie FABRIC_WORKSPACE_ID.

Important

Zmienna środowiskowa ustawiona z terminalu istnieje tylko wewnątrz tej jednej sesji terminalu. Nie jest współdzielony z innymi oknami terminala, z oknami używającymi innego typu powłoki ani z procesami uruchomionymi poza tym terminalem — w tym ze skryptami uruchamianymi za pomocą przycisku Uruchom w VS Code, który często uruchamia własny terminal. Jeśli skrypt nie może odnaleźć zmiennej, kończy się niepowodzeniem z parametrem Set FABRIC_WORKSPACE_ID environment variable before running the script.

Aby tego uniknąć, uruchom skrypt z sesji terminalu same gdzie ustawisz zmienną lub ustawisz ją trwale (zobacz sekcje Windows i macOS/Linux), aby każda nowa sesja terminalu pobierała ją automatycznie.

Ustawianie zmiennej środowiskowej na Windows

W systemie Windows zmienną można ustawić w dowolnym terminalu obsługującym zmienne środowiskowe — PowerShell, Windows PowerShell, oknach programu PowerShell lub Wiersza polecenia wbudowanych w programy Visual Studio i Visual Studio Code, Terminal Windows oraz w większości innych powłok.

Uruchom następujące polecenie w programie PowerShell lub zintegrowanym terminalu programu VS Code:

$env:FABRIC_WORKSPACE_ID="<WORKSPACE_ID>"

Aby potwierdzić, że zmienna jest ustawiona:

echo $env:FABRIC_WORKSPACE_ID

Spowoduje to ustawienie zmiennej tylko dla bieżącej sesji terminalu.

Ustawianie trwałej zmiennej środowiskowej (Windows)

Aby udostępnić zmienną w przyszłych sesjach, użyj jednej z następujących opcji:

  • PowerShell (jednolinijkowe polecenie): Uruchom setx FABRIC_WORKSPACE_ID "<WORKSPACE_ID>". Polecenie setx zapisuje dane w środowisku użytkownika, ale nie aktualizuje bieżącego terminalu — zamknij i otwórz ponownie terminal (lub otwórz nowy) przed uruchomieniem skryptu.
  • Graficzny interfejs użytkownika:
    1. Otwórz okno Właściwości systemu.
    2. Wybierz pozycję Zaawansowane ustawienia systemowe.
    3. Wybierz pozycję Zmienne środowiskowe.
    4. W sekcji Zmienne użytkownika wybierz pozycję Nowa.
    5. Wejść:
      • Nazwa: FABRIC_WORKSPACE_ID
      • Wartość: identyfikator obszaru roboczego
    6. Wybierz OK, aby zapisać.
    7. Zamknij i otwórz ponownie terminal przed ponownym uruchomieniem skryptu.

Ustawianie zmiennej środowiskowej w systemie macOS lub Linux

W systemach macOS i Linux można ustawić zmienną w dowolnej powłoce obsługującej export — Bash, Zsh (domyślna we współczesnych wersjach macOS), Fish (z nieco inną składnią), a także w zintegrowanych terminalach w Visual Studio Code i innych edytorach.

Run:

export FABRIC_WORKSPACE_ID="<WORKSPACE_ID>"

Aby potwierdzić, że zmienna jest ustawiona:

echo $FABRIC_WORKSPACE_ID

To ustawia zmienną tylko dla bieżącej sesji powłoki.

Ustawianie trwałej zmiennej środowiskowej (macOS lub Linux)

Aby udostępnić tę zmienną w przyszłych sesjach, dodaj wiersz export do profilu powłoki:

  • Zsh (ustawienie domyślne w systemie macOS): ~/.zshrc
  • Bash: ~/.bashrc (Linux) lub ~/.bash_profile (macOS)
  • Fish: uruchom set -Ux FABRIC_WORKSPACE_ID "<WORKSPACE_ID>" zamiast edytować plik

Po zaktualizowaniu profilu otwórz nowy terminal lub uruchom source ~/.zshrc (lub odpowiedni plik), aby zmiana obowiązywała.

Krok 4. Dodawanie funkcji pomocnika

W tym kroku wydzielasz zagadnienia przekrojowe — uwierzytelnianie, tworzenie nagłówków, odpytywanie operacji długotrwałych oraz logikę ponawiania prób — do niewielkiego zestawu funkcji pomocniczych wielokrotnego użytku, z których może korzystać każda funkcja kroku.

Centralizacja tych kwestii w funkcjach pomocniczych — zamiast umieszczania ich bezpośrednio w każdym miejscu wywołania — daje trzy konkretne korzyści:

  • Jedno źródło prawdy dla kwestii krzyżowych: Uwierzytelnianie, nagłówki i sondowanie LRO są potrzebne przez prawie każde wywołanie interfejsu API. Ich centralizacja sprawia, że każda funkcja krokowa koncentruje się na własnym zasobie, zamiast ponownie implementować mechanizmy pozyskiwania tokenów i ponawiania prób.
  • Odporność bez zbędnej złożoności: Funkcje pomocnicze przejmują obsługę stanów przejściowych — asynchronicznego aprowizowania, opóźnień propagacji po stronie backendu i przejściowych błędów wysyłania, które można ponowić — dzięki czemu funkcje krokowe pozostają krótkie i czytają się jak lista kontrolna.
  • Łatwiejsze do uczenia i modyfikowania: każdy pomocnik jest wprowadzany raz i ponownie używany. Jeśli Fabric zmieni wzorzec LRO lub zakres uwierzytelniania, poprawiasz to w jednym miejscu.

Pomocnicy dodani w tym kroku to:

  • Narzędzia pomocnicze do uwierzytelniania: tworzenie nagłówków dla interfejsów API REST usługi Fabric (oraz punktów końcowych LRO klastra usługi Power BI)
  • FabricClient: lekkie opakowanie do spójnych wywołań interfejsu API
  • Obsługa LRO: odpytywanie długotrwałych operacji za pomocą Location / x-ms-operation-id / Retry-After, w tym odpowiedzi typu 200-with-Running, punktów końcowych klastra Power BI i ładunków zakończenia zawierających wyłącznie stan (rozstrzygane przez displayName)
  • Pomocnik ładunku definicji: kodowanie map.json base64 dla definicji wbudowanych
  • Asystent połączenia Eventstream: monituje o podanie parametrów połączenia niestandardowego punktu końcowego
  • Pomocnik inicjatora: wysyła zdarzenia początkowe z logiką ponawiania prób w celu zapewnienia pomyślnego pozyskiwania
  • Pomocnik gotowości bazy danych KQL: czeka, aż baza danych KQL stanie się dostępna dla Fabric Maps

Note

Ten samouczek obejmuje dwie płaszczyzny:

  • Płaszczyzna sterowania (interfejsy API REST usługi Fabric): tworzenie zasobów Eventhouse, Eventstream i Map
  • Płaszczyzna danych/zapytań (interfejs API zarządzania Kusto): tworzenie tabel i funkcji KQL oraz zarządzanie nimi w obrębie eventhouse

Tworzenie funkcji pomocnika uwierzytelniania

Każde wywołanie interfejsu API REST usługi Fabric wykonywane w tym samouczku zawiera token dostępu Microsoft Entra (token okaziciela) w nagłówku Authorization. Zamiast pozyskiwać tokeny ad hoc, ten krok opakowuje DefaultAzureCredential w niewielki TokenProvider i udostępnia generator nagłówków specyficznych dla odbiorcy dla każdej rodziny punktów końcowych, którą wywołuje skrypt.

Scentralizowanie pozyskiwania tokenów i konstruowania nagłówków w pomocnikach — zamiast uzyskiwania tokenów w każdej lokacji połączeń — zapewnia trzy konkretne korzyści:

  • Scentralizowane poświadczenie: Pojedyncze DefaultAzureCredential jest opakowywane w TokenProvider i wykorzystywane ponownie przy każdym wywołaniu interfejsu API, więc wykrywanie tożsamości (Azure CLI, VS Code, tożsamość zarządzana itp.) odbywa się tylko raz.
  • Tokeny uwzględniające odbiorcę: Fabric, Kusto i punkty końcowe klastrów Power BI odrzucają tokeny wystawione dla niewłaściwego odbiorcy. Oddzielny konstruktor nagłówków dla każdej grupy odbiorców pozwala zachować właściwy zakres bezpośrednio obok miejsca wywołania, dzięki czemu od razu widać, do którego punktu końcowego jest kierowana każda funkcja.
  • Generowany przy każdym żądaniu: konstruktory nagłówków tworzą nagłówek Authorization na bieżąco, zamiast samodzielnie buforować token. Bazowe poświadczenie jest automatycznie odświeżane w tle, więc miejsca wywołania nigdy nie muszą martwić się jego wygaśnięciem.

W tym samouczku pokazano wywoływanie interfejsów API REST platformy Fabric przy użyciu zakresów delegowanych, takich jak Item.ReadWrite.All.

Dodaj następujące elementy po Config klasie:

# =========================================================
# Auth helpers
#
# Authentication utilities built on `DefaultAzureCredential` that acquire and
# construct Authorization headers for calling Fabric REST APIs.
# =========================================================

class TokenProvider:
    """
    Thin wrapper around `DefaultAzureCredential` that acquires Entra access
    tokens. `_fabric_headers()` and `_pbi_headers()` call `get()` per
    request so the Authorization header is always fresh; the underlying
    credential refreshes transparently.
    """
    def __init__(self):
        self._cred = DefaultAzureCredential()

    def get(self, scope: str) -> str:
        return self._cred.get_token(scope).token


_tokens = TokenProvider()


def _fabric_headers() -> dict[str, str]:
    """
    Build headers for Fabric REST API calls.

    This function is called each time we make a Fabric REST call so the token is fresh.
    """
    return {
        "Authorization": f"Bearer {_tokens.get('https://api.fabric.microsoft.com/.default')}",
        "Content-Type": "application/json"
    }


def _kusto_headers() -> dict[str, str]:
    """
    Build headers for Kusto (Eventhouse `queryServiceUri`) management and query calls.
    """
    return {
        "Authorization": f"Bearer {_tokens.get('https://api.kusto.windows.net/.default')}",
        "Content-Type": "application/json",
        "Accept": "application/json"
    }


def _pbi_headers() -> dict[str, str]:
    """
    Build headers for polling Power BI cluster LRO endpoints
    (e.g., df-*.analysis.windows.net) that require a Power BI audience token.
    """
    return {
        "Authorization": f"Bearer {_tokens.get('https://analysis.windows.net/powerbi/api/.default')}",
        "Content-Type": "application/json"
    }

Note

Niektóre długotrwałe operacje Fabric (LRO) są obsługiwane w punktach końcowych klastrów Power BI (*.analysis.windows.net), a nie w api.fabric.microsoft.com. Te punkty końcowe wymagają tokenu odbiorców Power BI, więc pomocnik LRO przełącza się na _pbi_headers() automatycznie po wykryciu tego adresu URL sondowania.

Utwórz opakowanie klienta Fabric

Większość wywołań REST Fabric w tym samouczku wysyła te same nagłówki Authorization i Content-Type. Zamiast powtarzać je w każdym miejscu wywołania, ten samouczek opakowuje httpx.Client w małą FabricClient, która automatycznie dołącza nagłówki, a jednocześnie nadal zwraca surowe httpx.Response, aby każde miejsce wywołania mogło sprawdzić kody stanu (na przykład, by odróżnić 201 od 202).

Opakowanie httpx.Client w ten sposób — zamiast przekazywania headers=_fabric_headers() w każdym miejscu wywołania — daje dwie konkretne korzyści:

  • Nagłówki w jednym miejscu: Każde miejsce wywołania automatycznie używa najnowszego _fabric_headers(), dzięki czemu nie da się przypadkowo wysłać nowego żądania bez nagłówka Authorization.
  • Kody stanu pozostają widoczne: request() zwraca surowe httpx.Response zamiast zdekodowanego formatu JSON, więc miejsca wywołania mogą nadal rozgałęziać logikę na podstawie kodu stanu (201 vs 202) i sprawdzać nagłówki, takie jak Location lub Retry-After, na potrzeby obsługi LRO.

Dodaj następujące elementy po funkcjach pomocnika uwierzytelniania:

# =========================================================
# FabricClient (minimal wrapper so call sites stay clean)
# =========================================================

class FabricClient:
    """
    Small wrapper around httpx.Client so we don't repeat headers everywhere.

    Keeps the tutorial behavior:
    - request() returns the raw httpx.Response so the caller can handle 201 vs 202.
    """
    def __init__(self, http_client: httpx.Client):
        self._http = http_client

    def request(self, method: str, url: str, *, json_body=None) -> httpx.Response:
        return self._http.request(method, url, headers=_fabric_headers(), json=json_body)

Utwórz funkcję pomocniczą LRO

Kilka interfejsów API REST Fabric używanych w tym samouczku — takich jak Tworzenie usługi Eventhouse, Tworzenie strumienia zdarzeń i Tworzenie mapy — obsługuje długotrwałe operacje (LROs).

Te interfejsy API mogą zwracać odpowiedzi w kilku wzorcach:

  • 201 Created z treścią zasobu wbudowaną (synchronicznie)
  • 202 Accepted z nagłówkiem Location wskazującym adres URL stanu operacji (asynchroniczny)
  • 202 Accepted z nagłówkiem x-ms-operation-id zamiast Location (asynchronicznej, alternatywnej formy)
  • 200 OK z status: "Running" lub status: "NotStarted" podczas sondowania (nadal w toku)
  • 200 OK z status: "Succeeded", ale bez identyfikatora zasobu w treści żądania (zakończono pomyślnie; rozwiąż problem, wyświetlając listę i dopasowując displayName)

Aby zapewnić spójne obsługę wszystkich tych elementów, należy utworzyć jedną funkcję pomocnika, która:

  1. Zwraca identyfikator zasobu natychmiast, jeśli początkowa odpowiedź już go zawiera.
  2. W przeciwnym razie sprawdza adres URL operacji (utworzony na podstawie Location lub x-ms-operation-id) za pomocą Retry-After.
  3. 200 OK status: "Running" / "NotStarted"Traktuje jako nadal w toku i kontynuuje sondowanie.
  4. W razie powodzenia zwraca identyfikator zasobu z treści odpowiedzi, a gdy treść odpowiedzi zawiera wyłącznie status, przechodzi do listowania zasobów i dopasowuje je na podstawie displayName (z ponownymi próbami).
  5. Używa _pbi_headers(), gdy adres URL odpytywania znajduje się w klastrze Power BI (*.analysis.windows.net), a w przeciwnym razie używa nagłówków Fabric.

Ta jedna funkcja pomocnicza eliminuje potrzebę stosowania osobnych funkcji „wyszukiwania po nazwie” dla każdego zasobu — każda funkcja create_* w tym samouczku wywołuje _handle_lro z odpowiednimi list_url i match_display_name.

Dodaj następujące elementy po FabricClient klasie:

# =========================================================
# LRO handler 
# =========================================================

def _handle_lro(
    client: httpx.Client,
    initial_response: httpx.Response,
    *,
    list_url: str | None = None,
    match_display_name: str | None = None,
    id_field: str = "id",
    max_attempts: int = 10,
    delay: int = 5,
) -> str:
    """
    Handle a Fabric long-running operation (LRO) and return the resource id.

    Supports the response patterns used by Fabric REST APIs:
    - 200/201 with the resource body inline (synchronous).
    - 202 with a `Location` header or `x-ms-operation-id` (asynchronous).
    - 200 with `status: "Running"` / `"NotStarted"` while polling.
    - 200 with `status: "Succeeded"` but no id (resolve by listing and matching `displayName`).

    Polling uses `Retry-After` and switches to a Power BI audience token when
    the operation URL is on `*.analysis.windows.net`.
    """
    # Sync 200/201 with body: return the id immediately.
    if initial_response.status_code in (200, 201):
        try:
            body = initial_response.json() if initial_response.content else {}
        except ValueError:
            body = {}
        if isinstance(body, dict) and body.get(id_field):
            return body[id_field]

    # Location header, with x-ms-operation-id fallback.
    op_url = initial_response.headers.get("Location")
    if not op_url:
        op_id = initial_response.headers.get("x-ms-operation-id")
        if op_id:
            op_url = f"https://api.fabric.microsoft.com/v1/operations/{op_id}"
        else:
            raise RuntimeError(
                f"Missing LRO Location/x-ms-operation-id. "
                f"status={initial_response.status_code} body={initial_response.text[:500]!r}"
            )

    # Audience-aware polling: Power BI cluster endpoints need a different token.
    poll_headers = _pbi_headers() if "analysis.windows.net" in op_url else _fabric_headers()
    retry_after = int(initial_response.headers.get("Retry-After", "5"))

    while True:
        time.sleep(retry_after)
        poll = client.get(op_url, headers=poll_headers)

        if poll.status_code == 202:
            retry_after = int(poll.headers.get("Retry-After", "5"))
            continue

        poll.raise_for_status()
        body = poll.json() if poll.content else {}
        status = body.get("status") if isinstance(body, dict) else None

        if status in ("Running", "NotStarted"):
            retry_after = int(poll.headers.get("Retry-After", "5"))
            continue
        if status == "Failed":
            raise RuntimeError(f"LRO failed. Body: {body}")

        if isinstance(body, dict) and body.get(id_field):
            return body[id_field]

        # Status-only success: list and match by displayName, with retries.
        if status == "Succeeded" and list_url and match_display_name:
            for attempt in range(max_attempts):
                r = client.get(list_url, headers=_fabric_headers())
                r.raise_for_status()
                match = next(
                    (i for i in r.json().get("value", []) if i.get("displayName") == match_display_name),
                    None,
                )
                if match and match.get(id_field):
                    return match[id_field]
                time.sleep(delay)
            raise RuntimeError(
                f"LRO succeeded but resource not visible after retries. "
                f"match_display_name={match_display_name!r}"
            )

        raise RuntimeError(f"LRO completed but no resource id was returned. Body: {body}")

Note

Nowo utworzone zasoby mogą nie być widoczne od razu podczas wywoływania interfejsów API służących do listowania z powodu opóźnień propagacji po stronie zaplecza. Funkcja pomocnika automatycznie ponawia próbę, dopóki zasób nie stanie się widoczny.

Pomocnik ładunku definicji

Podczas tworzenia mapy z definicją publiczną interfejs API REST Create map oczekuje, że każda część w definition.parts będzie zawierać ładunek zakodowany w formacie base64 przy użyciu "payloadType": "InlineBase64". Funkcja pomocnicza _json_to_b64 koduje obiekt Python dict (Twój map.json) do tego formatu, aby create_map mogło wstawić go bezpośrednio do treści żądania.

Dodaj następujące elementy po _handle_lro funkcji:

# =========================================================
# Definition payload helper
#
# Encodes map.json as base64 for inline Create map payloads.
# =========================================================

def _json_to_b64(obj: dict) -> str:
    """
    Convert a Python dict to base64-encoded JSON text.

    Fabric Map "Create Map with definition inline" requires:
    - definition.parts[].payloadType = InlineBase64
    - definition.parts[].payload     = base64(json(map_json))
    """
    return base64.b64encode(json.dumps(obj).encode("utf-8")).decode("utf-8")

Utworzenie pomocnika udostępniającego parametry połączenia dla strumienia zdarzeń

Aby wysyłać zdarzenia do niestandardowego punktu końcowego strumienia zdarzeń, skrypt wymaga parametry połączenia dla tego punktu końcowego.

W przeciwieństwie do interfejsów API REST usługi Fabric, które były dotąd wywoływane (czyli operacji płaszczyzny sterowania służących do tworzenia zasobów i zarządzania nimi), pozyskiwanie danych strumienia zdarzeń używa punktu końcowego płaszczyzny danych zgodnego z usługą Event Hubs, a ten punkt końcowy uwierzytelnia się przy użyciu parametrów połączenia opartych na SAS, a nie tokenu Microsoft Entra. Parametry połączenia są tworzone podczas dodawania niestandardowego źródła punktu końcowego i nie są udostępniane przez interfejs API REST usługi Fabric, dlatego trzeba je skopiować z portalu Fabric.

get_eventhub_connection_string_interactive albo używa wartości ze zmiennej środowiskowej EVENTHUB_CONNECTION_STRING (co jest przydatne przy ponownych uruchomieniach), albo prosi o jej podanie w czasie wykonywania, a następnie zapisuje ją w pamięci podręcznej na cfg, aby kolejne kroki mogły użyć jej ponownie bez ponawiania monitu.

Dodaj następujące elementy po _json_to_b64 funkcji:

def get_eventhub_connection_string_interactive(cfg: Config) -> str:
    """
    Prompt for (or read) the eventstream custom endpoint connection string.

    The connection string is created when the custom endpoint source is added
    to the eventstream and isn't exposed by the Fabric REST API, so we read it
    from the `EVENTHUB_CONNECTION_STRING` environment variable when set, or
    prompt interactively otherwise. The value is cached on `cfg` for reuse.
    """
    if getattr(cfg, "eventhub_connection_string", None):
        return cfg.eventhub_connection_string

    print("\n=== Eventstream connection string required ===")
    print("In the Fabric portal:")
    print("  1) Open the eventstream you just created")
    print("  2) Select the custom endpoint source")
    print("  3) Select SAS Key Authentication")
    print("  4) Copy Connection string-primary key\n")

    cfg.eventhub_connection_string = input("Paste connection string here: ").strip()

    if not cfg.eventhub_connection_string:
        raise RuntimeError("Connection string cannot be empty.")

    return cfg.eventhub_connection_string

Note

Parametry połączenia są oddzielne od tokenu dostępu Microsoft Entra używanego przez interfejsy API REST usługi Fabric. Token interfejsu API REST służy do zarządzania zasobami, natomiast parametr połączenia strumienia zdarzeń służy do pozyskiwania danych w trybie strumieniowym.

Tworzenie pomocnika w celu inicjowania zdarzeń początkowych z pliku CSV

Aby upewnić się, że mapa wyświetla dane natychmiast po jej utworzeniu, skrypt wysyła niewielki zestaw zdarzeń początkowych do strumienia zdarzeń przed utworzeniem mapy.

Bez tego kroku tabela Eventhouse może jeszcze nie zawierać żadnych danych, a mapa może być pusta podczas pierwszego ładowania.

Ta funkcja pomocnika odczytuje dane z lokalnego pliku CSV i wysyła każdy wiersz jako zdarzenie JSON do strumienia zdarzeń przy użyciu protokołu EventHub.

Ponieważ zasoby Eventstream są aprowizowane asynchronicznie, niestandardowy punkt końcowy może nie być gotowy do przyjmowania zdarzeń bezpośrednio po utworzeniu. Aby sobie z tym poradzić, funkcja pomocnicza zawiera wbudowaną logikę ponawiania prób, która automatycznie podejmuje próby wysyłania zdarzeń do czasu, aż punkt końcowy będzie dostępny. Dzięki temu proces inicjowania jest niezawodny i powtarzalny i nie wymaga ręcznych korekt chronometrażu.

To podejście odzwierciedla rzeczywiste wzorce pozyskiwania danych:

  • Dane są tworzone zewnętrznie (na przykład urządzenia IoT lub aplikacje)
  • Zdarzenia są strumieniowane do Eventstream
  • Eventstream dostarcza dane do Eventhouse do wykonywania zapytań i wizualizacji

Dodając zdarzenia początkowe, symulujesz ten przepływ pozyskiwania danych i zapewniasz, że:

  • Tabela docelowa jest wypełniana
  • Funkcja KQL ma dane do zwrócenia
  • Mapa jest renderowana natychmiast po utworzeniu

Dodaj następujący kod po get_eventhub_connection_string_interactive() funkcji:

def seed_eventstream_from_csv(cfg: Config, max_attempts: int = 10, delay: int = 3) -> int:
    """
    Send seed events from a CSV with retries to handle eventstream readiness delay.
    """
    conn_str = get_eventhub_connection_string_interactive(cfg)

    last_error = None
    for attempt in range(1, max_attempts + 1):
        print(f"Seeding attempt {attempt}/{max_attempts}...")
        try:
            sent = 0
            producer = EventHubProducerClient.from_connection_string(conn_str=conn_str)
            try:
                with open(cfg.seed_csv_path, newline="", encoding="utf-8") as f:
                    reader = csv.DictReader(f)
                    batch = producer.create_batch()
                    batch_count = 0
                    for row in reader:
                        event = {
                            "VehicleId": row["VehicleId"],
                            "Latitude": float(row["Latitude"]),
                            "Longitude": float(row["Longitude"]),
                            "EventTime": row["EventTime"],
                        }
                        data = EventData(json.dumps(event))
                        try:
                            batch.add(data)
                            batch_count += 1
                        except ValueError:
                            producer.send_batch(batch)
                            sent += batch_count
                            batch = producer.create_batch()
                            batch.add(data)
                            batch_count = 1
                    if batch_count > 0:
                        producer.send_batch(batch)
                        sent += batch_count
                print(f"Seed events sent: {sent}")
                return sent
            finally:
                producer.close()
        except (EventHubError, ConnectionError, TimeoutError) as exc:
            last_error = exc
            print(f"Seeding failed (attempt {attempt}): {exc}. Retrying in {delay}s...")
            time.sleep(delay)

    raise RuntimeError(f"Seeding failed after {max_attempts} attempts. Last error: {last_error}")

Note

Ten krok wprowadza niewielką ilość danych statycznych do potoku przesyłania strumieniowego.
W scenariuszu produkcyjnym zdarzenia są zwykle generowane w sposób ciągły przez systemy zewnętrzne, a nie ładowane z pliku.

Czekaj na dostępność bazy danych KQL

Gdy magazyn zdarzeń, jego baza danych KQL i funkcja KQL istnieją, baza danych KQL może nadal nie być natychmiast rozpoznawana z innych punktów końcowych REST Fabric. Usługi Fabric działają w rozproszonych systemach zaplecza, więc propagacja nowo utworzonego zasobu może zająć krótką chwilę.

Jeśli wywołasz metodę Create Map bezpośrednio po utworzeniu funkcji KQL, tworzenie mapy może nie rozpoznać źródła danych i zwrócić błąd, taki jak nie znaleziono bazy danych Kusto.

wait_for_kql_database_ready odpytuje punkt końcowy REST usługi Fabric dla bazy danych KQL i zwraca wynik, gdy tylko otrzyma odpowiedź 200 OK. Jest to mechanizm kontroli działający w trybie best effort — pomyślna odpowiedź z tego punktu końcowego płaszczyzny sterowania stanowi mocną przesłankę, że usługa Maps może również uzyskać dostęp do bazy danych — i zgłasza RuntimeError po max_attempts, jeśli baza danych nigdy nie stanie się widoczna.

Dodaj następujące elementy po seed_eventstream_from_csv funkcji:

# =========================================================
# KQL database readiness helper
#
# Polls the KQL database's Fabric REST endpoint until it
# responds 200, as a best-effort gate before Create Map
# references it as a data source.
# =========================================================

def wait_for_kql_database_ready(
    client: httpx.Client,
    cfg: Config,
    kql_database_item_id: str,
    max_attempts: int = 10,
    delay: int = 3,
) -> None:
    """
    Poll the Fabric REST endpoint for a KQL database until it returns 200.

    Acts as a best-effort readiness gate before calling Create Map with the
    KQL database as a data source. Retries `max_attempts` times with `delay`
    seconds between attempts, then raises `RuntimeError` if the database
    never becomes visible.
    """
    url = f"https://api.fabric.microsoft.com/v1/workspaces/{cfg.workspace_id}/kqlDatabases/{kql_database_item_id}"
    for attempt in range(1, max_attempts + 1):
        resp = client.get(url, headers=_fabric_headers())
        if resp.status_code == 200:
            print("KQL database is available to Fabric Maps")
            return
        print(
            f"Waiting for KQL database availability "
            f"(attempt {attempt}/{max_attempts}, status={resp.status_code})..."
        )
        time.sleep(delay)

    raise RuntimeError(
        f"KQL database {kql_database_item_id!r} did not become available after {max_attempts} attempts."
    )

Tworzenie podstawowych funkcji

Następnie dodasz podstawowe funkcje, które definiują przepływ pracy. Wszystkie są wywoływane z main().

Funkcje są dodawane w kolejności, w której są zdefiniowane w kodzie. main() wywołuje je w nieco innej kolejności, tak aby tabela KQL istniała, zanim strumień zdarzeń zostanie z nią powiązany, i aby dane inicjalizacyjne były dostępne, zanim zostanie przeprowadzona weryfikacja.

  • Stworzyć dom eventowy
  • Tworzenie tabeli KQL
  • Weryfikowanie pozyskiwania (wywoływane po inicjowaniu)
  • Tworzenie strumienia zdarzeń
  • Tworzenie funkcji KQL
  • Tworzenie definicji mapy (map.json)
  • Tworzenie metadanych platformy (.platform)
  • Tworzenie mapy

Stworzyć dom eventowy

create_eventhouse tworzy element Eventhouse w obszarze roboczym i zwraca jego identyfikator. Interfejs API REST tworzenia usługi Eventhouse może odpowiedzieć na to samo wywołanie na trzy różne sposoby:

  • 201 Created z identyfikatorem eventhouse w tekście (synchronicznie).
  • 202 Accepted z adresem URL operacji LRO (asynchronicznym).
  • 409 Conflict z ustawieniem x-ms-public-api-error-code na ItemDisplayNameNotAvailableYet (poprzednia nazwa jest nadal zarezerwowana po stronie backendu) lub ItemDisplayNameAlreadyInUse (eventhouse o tej nazwie już istnieje w obszarze roboczym).

Aby niezawodnie obsłużyć wszystkie trzy elementy: create_eventhouse

  • Deleguje odpowiedzi 201 i 202 do _handle_lro, które już obsługuje synchroniczne i asynchroniczne zakończenie w jednolity sposób.
  • Uwzględnia Retry-After i ponawia próby (maksymalnie pięć razy) na ItemDisplayNameNotAvailableYet.
  • Ponownie używa istniejącego eventhouse w elemencie ItemDisplayNameAlreadyInUse, wyświetlając listę eventhouse'ów w obszarze roboczym i dopasowując według displayName.
  • W ostateczności używa unikalnej nazwy wyświetlanej (z krótkim sufiksem UUID), jeśli nazwa nigdy nie stanie się dostępna po wyczerpaniu limitu ponowień, dzięki czemu skrypt może nadal kontynuować działanie.

Dodaj następujące elementy po wait_for_kql_database_ready funkcji:

# =========================================================
# Create an eventhouse
# =========================================================

def create_eventhouse(client: httpx.Client, fabric: FabricClient, cfg: Config) -> str:
    """
    Create an eventhouse in the workspace and return its item ID.

    Handles the three response patterns Create Eventhouse can return:
    - 201/202: delegate to `_handle_lro` (synchronous body or LRO completion).
    - 409 `ItemDisplayNameNotAvailableYet`: honor `Retry-After` and retry.
    - 409 `ItemDisplayNameAlreadyInUse`: list eventhouses and reuse the one
      whose `displayName` matches `cfg.eventhouse_display_name`.

    If the name remains unavailable after the retry budget, falls back to a
    uniquified display name (suffixed with a short UUID) so the script can
    still make forward progress.
    """
    eventhouse_url = f"https://api.fabric.microsoft.com/v1/workspaces/{cfg.workspace_id}/eventhouses"
    eventhouse_payload = {
        "displayName": cfg.eventhouse_display_name,
        "description": cfg.eventhouse_description
    }

    # Retry loop to handle transient "name not available yet"
    for attempt in range(1, 6):
        eh_resp = fabric.request("POST", eventhouse_url, json_body=eventhouse_payload)

        print("Create eventhouse status:", eh_resp.status_code)
        print("Create eventhouse headers:", dict(eh_resp.headers))

        # 201/202: success or LRO — _handle_lro handles both.
        if eh_resp.status_code in (201, 202):
            eventhouse_id = _handle_lro(
                client, eh_resp,
                list_url=eventhouse_url,
                match_display_name=cfg.eventhouse_display_name,
            )
            print("Eventhouse created. Eventhouse ID:", eventhouse_id)
            return eventhouse_id

        # 409: name issues
        if eh_resp.status_code == 409:
            api_code = eh_resp.headers.get("x-ms-public-api-error-code")

            # Name reserved temporarily: wait and retry
            if api_code == "ItemDisplayNameNotAvailableYet":
                wait_s = int(eh_resp.headers.get("retry-after", "20"))
                print(f"Name not available yet (attempt {attempt}/5). Waiting {wait_s}s then retrying...")
                time.sleep(wait_s)
                continue

            # Name already exists: reuse existing eventhouse by displayName
            if api_code == "ItemDisplayNameAlreadyInUse":
                print(f"Eventhouse {cfg.eventhouse_display_name!r} already exists. Reusing it...")
                r = client.get(eventhouse_url, headers=_fabric_headers())
                r.raise_for_status()
                match = next(
                    (i for i in r.json().get("value", []) if i.get("displayName") == cfg.eventhouse_display_name),
                    None,
                )
                if match and match.get("id"):
                    return match["id"]
                raise RuntimeError(
                    f"Eventhouse {cfg.eventhouse_display_name!r} reported as existing but not found in list."
                )

        # Anything else: fail fast with details
        raise RuntimeError(f"Create eventhouse failed: {eh_resp.status_code} {eh_resp.text}")

    # If the name never becomes available, last-resort: pick a unique name and try once
    cfg.eventhouse_display_name = f"{cfg.eventhouse_display_name}-{uuid.uuid4().hex[:8]}"
    print(f"Name still not available; switching to unique name: {cfg.eventhouse_display_name}")
    eh_resp = fabric.request("POST", eventhouse_url, json_body={
        "displayName": cfg.eventhouse_display_name,
        "description": cfg.eventhouse_description
    })
    eh_resp.raise_for_status()
    return eh_resp.json()["id"]

Tworzenie tabeli KQL

create_kql_table_if_missing gwarantuje, że tabela docelowa istnieje w bazie danych KQL przed rozpoczęciem zapisywania w nim strumienia zdarzeń. Miejsce docelowe strumienia zdarzeń, które utworzysz później, jest skonfigurowane z użyciem ProcessedIngestion i stałego tableName, więc tabela musi już istnieć, gdy zdarzenia zaczną napływać — w przeciwnym razie pozyskiwanie danych zakończy się niepowodzeniem.

Funkcja wysyła polecenie .create-merge table do punktu końcowego zarządzania Kusto usługi Eventhouse (queryServiceUri + /v1/rest/mgmt). .create-merge jest idempotentne: tworzy tabelę, jeśli ona nie istnieje, a jeśli już istnieje, scala schemat. Dzięki temu można ją bezpiecznie wywoływać przy każdym uruchomieniu.

Przed wydaniem polecenia funkcja odczytuje właściwości eventhouse, aby uzyskać queryServiceUri i identyfikator elementu bazy danych KQL, a następnie ustala wartość displayName bazy danych, tak aby ładunek mgmt odwoływał się do niej według nazwy, a nie identyfikatora.

Dodaj następujące elementy po create_eventhouse funkcji:

# =========================================================
# Create the KQL table (idempotent)
# =========================================================

def create_kql_table_if_missing(client: httpx.Client, fabric: FabricClient, cfg: Config, eventhouse_id: str) -> None:
    """
    Create or merge the destination table in the eventhouse's KQL database.

    Reads the eventhouse properties to discover `queryServiceUri` and the KQL
    database item ID, resolves the database's `displayName`, then issues a
    `.create-merge table` command against the Kusto management endpoint.
    `.create-merge` is idempotent: it creates the table if missing and merges
    the schema if it already exists.
    """
    # Get queryServiceUri + KQL database item id
    get_eventhouse_url = f"https://api.fabric.microsoft.com/v1/workspaces/{cfg.workspace_id}/eventhouses/{eventhouse_id}"
    eh = fabric.request("GET", get_eventhouse_url)
    eh.raise_for_status()

    props = eh.json().get("properties") or {}
    query_service_uri = props.get("queryServiceUri")
    databases_item_ids = props.get("databasesItemIds") or []
    if not query_service_uri or not databases_item_ids:
        raise RuntimeError("Eventhouse missing queryServiceUri or databasesItemIds")

    kql_database_item_id = databases_item_ids[0]

    # Resolve actual DB displayName (don't rely on cfg.kql_database_name)
    get_db_url = f"https://api.fabric.microsoft.com/v1/workspaces/{cfg.workspace_id}/kqlDatabases/{kql_database_item_id}"
    db_resp = fabric.request("GET", get_db_url)
    db_resp.raise_for_status()
    kql_database_name = db_resp.json().get("displayName")
    if not kql_database_name:
        raise RuntimeError("KQL database response did not include displayName")

    # Create (or merge) the table schema
    # (Schema matches what your CSV sends: VehicleId, Latitude, Longitude, EventTime)
    csl = f""".create-merge table {cfg.eventhouse_table_name} (
        VehicleId: string,
        Latitude: real,
        Longitude: real,
        EventTime: datetime
    )"""

    mgmt_url = f"{query_service_uri}/v1/rest/mgmt"
    mgmt_payload = {"db": kql_database_name, "csl": csl}
    resp = client.post(mgmt_url, headers=_kusto_headers(), json=mgmt_payload)
    if resp.status_code >= 400:
        raise RuntimeError(f"Create table failed: {resp.status_code}\n{resp.text}")

    print(f"Ensured table exists: {cfg.eventhouse_table_name}")

Weryfikowanie pozyskiwania danych

verify_eventhouse_data potwierdza, że zdarzenia umieszczone w strumieniu zdarzeń faktycznie trafiły do tabeli eventhouse. Cyklicznie wysyła zapytanie <table> | count do punktu końcowego Kusto query usługi Eventhouse (queryServiceUri + /v1/rest/query), aż zwrócona liczba będzie większa od zera, albo kończy się błędem po upływie limitu czasu. Przesłanie danych z Eventstream z niestandardowego punktu końcowego do tabeli trwa kilka sekund, dlatego to cykliczne odpytywanie — a nie pojedyncze zapytanie — daje wiarygodny sygnał powodzenia lub niepowodzenia.

Jest zdefiniowany obok create_kql_table_if_missing, ponieważ obie funkcje pomocnicze odwołują się do tych samych właściwości magazynu zdarzeń (queryServiceUri, databasesItemIds) i ustalają wartość displayName bazy danych KQL. Jest wywoływana z main()poseed_eventstream_from_csv, więc zasiane zdarzenia mają szansę przepłynąć przez strumień zdarzeń i dotrzeć do tabeli, zanim zostanie wykonane zliczanie.

Uruchomienie tej kontroli z wyprzedzeniem pozwala wcześnie wykryć błędną konfigurację przetwarzania danych — na przykład miejsce docelowe strumienia zdarzeń skierowane do tabeli o nieprawidłowej nazwie — zamiast dopuścić, by problem ujawnił się później w postaci pustej mapy.

Dodaj następujące elementy po create_kql_table_if_missing funkcji:

# =========================================================
# Verify data ingestion (called after seeding)
# =========================================================

def verify_eventhouse_data(client: httpx.Client, fabric: FabricClient, cfg: Config, eventhouse_id: str):
    """
    Poll a count query against the eventhouse table until rows arrive.

    Reads the eventhouse properties to get `queryServiceUri` and the KQL
    database item ID, resolves the database's `displayName`, then polls
    `<table> | count` against the Kusto query endpoint
    (`queryServiceUri` + `/v1/rest/query`) until the count is greater than
    zero or the timeout elapses. Eventstream ingestion is asynchronous, so
    polling avoids a false negative when the query runs before seeded
    events have landed in the table.
    """

    # Reuse your existing pattern to get KQL DB info
    eh = fabric.request("GET", f"https://api.fabric.microsoft.com/v1/workspaces/{cfg.workspace_id}/eventhouses/{eventhouse_id}")
    eh.raise_for_status()

    props = eh.json().get("properties") or {}
    db_ids = props.get("databasesItemIds") or []
    query_service_uri = props.get("queryServiceUri")

    if not db_ids or not query_service_uri:
        raise RuntimeError("Missing eventhouse properties for verification")

    db_id = db_ids[0]

    db_resp = fabric.request("GET", f"https://api.fabric.microsoft.com/v1/workspaces/{cfg.workspace_id}/kqlDatabases/{db_id}")
    db_resp.raise_for_status()

    db_name = db_resp.json().get("displayName")

    # Simple count query
    csl = f"{cfg.eventhouse_table_name} | count"
    query_url = f"{query_service_uri}/v1/rest/query"

    max_attempts = 12
    delay_seconds = 5

    for attempt in range(1, max_attempts + 1):
        resp = client.post(
            query_url,
            headers=_kusto_headers(),
            json={"db": db_name, "csl": csl}
        )

        if resp.status_code >= 400:
            raise RuntimeError(f"Verification query failed: {resp.text}")

        # Kusto v1 query response: Tables[0].Rows[0][0] holds the count.
        count = resp.json()["Tables"][0]["Rows"][0][0]
        print(f"Data verification attempt {attempt}/{max_attempts}: count = {count}")

        if count > 0:
            print(f"Data ingestion verified: {count} row(s) in {cfg.eventhouse_table_name}")
            return

        if attempt < max_attempts:
            time.sleep(delay_seconds)

    raise RuntimeError(
        f"Data verification failed: no rows in {cfg.eventhouse_table_name} after "
        f"{max_attempts * delay_seconds}s"
    )

Tworzenie strumienia zdarzeń z definicją

create_eventstream_with_definition Tworzy strumień zdarzeń w obszarze roboczym z pełną topologią upieczętowaną w żądaniu, a następnie zwraca identyfikator elementu strumienia zdarzeń. Użycie publicznej definicji umożliwia utworzenie strumienia zdarzeń i skonfigurowanie jego źródeł, strumieni oraz miejsc docelowych w jednym wywołaniu, zamiast najpierw tworzyć strumień zdarzeń, a następnie modyfikować jego definicję.

Przed wysłaniem żądania funkcja odczytuje właściwości Eventhouse, aby uzyskać identyfikator elementu bazy danych KQL, i ustala displayName bazy danych, aby obiekt docelowy odwoływał się do niej według nazwy, a nie identyfikatora. Następnie tworzy graf strumienia zdarzeń ze źródłem CustomEndpoint, elementem DefaultStream oraz miejscem docelowym Eventhouse skonfigurowanym za pomocą ProcessedIngestion i stałej wartości tableName z cfg, koduje graf w formacie Base64 jako część eventstream.json i wysyła go metodą POST do Create Eventstream.

Interfejs API REST tworzenia strumienia zdarzeń może odpowiedzieć za pomocą 201 Created (synchronicznego, wbudowanego treści), 202 Accepted (asynchronicznego LRO za pośrednictwem Location lub x-ms-operation-id), lub 200 OK z ładunkiem uzupełniania tylko do stanu, gdzie strumień zdarzeń nie jest jeszcze widoczny w odpowiedzi List Eventstreams z powodu opóźnienia propagacji zaplecza. _handle_lro obejmuje wszystkie te przypadki — w tym tworzenie listy i dopasowywanie według displayName — dlatego ta funkcja przekazuje jej pełną obsługę odpowiedzi w jednym wywołaniu.

Dodaj następujące elementy po verify_eventhouse_data funkcji:

# =========================================================
# Create eventstream with definition
# =========================================================

def create_eventstream_with_definition(client: httpx.Client, fabric: FabricClient, cfg: Config, eventhouse_id: str) -> str:
    """
    Create an eventstream with a public definition and return its item ID.

    Reads the eventhouse properties to discover the KQL database item ID and
    resolves the database's `displayName`, then builds an eventstream graph
    with a `CustomEndpoint` source, a `DefaultStream`, and an `Eventhouse`
    destination configured with `ProcessedIngestion` and the table name from
    `cfg`. Base64-encodes the graph as the `eventstream.json` part, POSTs it
    to Create Eventstream, and delegates response handling to `_handle_lro`.
    """
    eventstream_url = f"https://api.fabric.microsoft.com/v1/workspaces/{cfg.workspace_id}/eventstreams"

    source_name = "CustomEndpointSource"
    stream_name = "DefaultStream"
    destination_name = "EventhouseDestination"

    # Resolve the KQL database item ID from the Eventhouse
    get_eventhouse_url = f"https://api.fabric.microsoft.com/v1/workspaces/{cfg.workspace_id}/eventhouses/{eventhouse_id}"
    eh = fabric.request("GET", get_eventhouse_url)
    eh.raise_for_status()

    props = (eh.json().get("properties") or {})
    databases_item_ids = props.get("databasesItemIds") or []
    if not databases_item_ids:
        raise RuntimeError("Eventhouse properties did not include databasesItemIds.")

    kql_database_item_id = databases_item_ids[0]

    # Resolve the actual KQL database *name* (displayName) to avoid name drift
    get_db_url = f"https://api.fabric.microsoft.com/v1/workspaces/{cfg.workspace_id}/kqlDatabases/{kql_database_item_id}"
    db_resp = fabric.request("GET", get_db_url)
    db_resp.raise_for_status()
    kql_database_name = db_resp.json().get("displayName")
    if not kql_database_name:
        raise RuntimeError("KQL database response did not include displayName.")

    print("Eventhouse ID:", eventhouse_id)
    print("KQL database item ID:", kql_database_item_id)
    print("KQL database name:", kql_database_name)

    eventstream_json = {
        "sources": [
            {
                "name": source_name,
                "type": "CustomEndpoint",
                "properties": {
                    "inputSerialization": {"type": "Json", "properties": {"encoding": "UTF8"}}
                }
            }
        ],
        "streams": [
            {
                "name": stream_name,
                "type": "DefaultStream",
                "properties": {},
                "inputNodes": [{"name": source_name}]
            }
        ],
        "operators": [],
        "destinations": [
            {
                "name": destination_name,
                "type": "Eventhouse",
                "properties": {
                    "dataIngestionMode": "ProcessedIngestion",
                    "workspaceId": cfg.workspace_id,
                    "itemId": kql_database_item_id,
                    "databaseName": kql_database_name,
                    "tableName": cfg.eventhouse_table_name,
                    "inputSerialization": {"type": "Json", "properties": {"encoding": "UTF8"}}
                },
                "inputNodes": [{"name": stream_name}]
            }
        ],
        "compatibilityLevel": "1.1"
    }

    eventstream_payload = {
        "displayName": cfg.eventstream_display_name,
        "description": cfg.eventstream_description,
        "definition": {
            "parts": [
                {
                    "path": "eventstream.json",
                    "payload": _json_to_b64(eventstream_json),
                    "payloadType": "InlineBase64"
                }
            ]
        }
    }

    es_resp = fabric.request("POST", eventstream_url, json_body=eventstream_payload)

    eventstream_id = _handle_lro(
        client,
        es_resp,
        list_url=eventstream_url,
        match_display_name=cfg.eventstream_display_name,
    )

    print("Eventstream created. Eventstream ID:", eventstream_id)
    return eventstream_id

Utwórz funkcję KQL

create_kql_function tworzy (lub aktualizuje) zapisaną funkcję Kusto w bazie danych KQL eventhouse i zwraca identyfikator elementu tej bazy danych KQL, aby wywołujący mógł połączyć z nią źródło danych mapy. Funkcja — LatestVehicleLocations domyślnie — zwraca najnowszy wiersz dla każdego VehicleId za pośrednictwem arg_max(EventTime, *), rzutując Latitude, Longitude, VehicleId i EventTime, aby usługa Fabric Maps mogła powiązać kolumny szerokości i długości geograficznej warstwy.

Podobnie jak create_kql_table_if_missing, to narzędzie pomocnicze działa względem punktu końcowego zarządzania Kusto usługi Eventhouse (queryServiceUri + /v1/rest/mgmt) i jest idempotentne: .create-or-alter function tworzy funkcję, jeśli jeszcze nie istnieje, a jeśli istnieje, zastępuje jej definicję, więc można je bezpiecznie wywoływać przy każdym uruchomieniu.

Polecenie jest wysyłane z użyciem skipvalidation=true, ponieważ treść funkcji odwołuje się do tabeli docelowej za pośrednictwem table("<name>"), zamiast jako sam identyfikator. Postać table() odracza rozpoznawanie nazw do czasu wykonywania zapytania, więc w przeciwnym razie walidacja w momencie tworzenia zakończyłaby się niepowodzeniem, jeśli tabela nie otrzymała jeszcze żadnych danych, a jej schemat nie jest jeszcze w pełni widoczny dla walidatora. Parowanie skipvalidation=true z funkcją table("...") umożliwia utworzenie funkcji przed wypełnieniem tabeli, czyli kolejności działania tego samouczka.

Dodaj następujące elementy po create_eventstream_with_definition funkcji:

# =========================================================
# Create KQL function
# =========================================================

def create_kql_function(client: httpx.Client, fabric: FabricClient, cfg: Config, eventhouse_id: str) -> str:
    """
    Create or update the stored Kusto function used by the map layer.

    Reads the eventhouse properties to discover `queryServiceUri` and the
    KQL database item ID, resolves the database's `displayName`, then
    issues a `.create-or-alter function` command against the Kusto
    management endpoint with `skipvalidation=true` and a `table("...")`
    reference so the function can be created before the destination table
    has any data. Returns the KQL database item ID so the caller can wire
    the map's data source to it.
    """
    # Get eventhouse properties (queryServiceUri + databasesItemIds)
    get_eventhouse_url = f"https://api.fabric.microsoft.com/v1/workspaces/{cfg.workspace_id}/eventhouses/{eventhouse_id}"
    eh = fabric.request("GET", get_eventhouse_url)
    eh.raise_for_status()

    props = (eh.json().get("properties") or {})
    query_service_uri = props.get("queryServiceUri")
    databases_item_ids = props.get("databasesItemIds") or []

    if not query_service_uri:
        raise RuntimeError("Eventhouse properties did not include queryServiceUri.")
    if not databases_item_ids:
        raise RuntimeError("Eventhouse properties did not include databasesItemIds.")

    # We'll return this so the caller can wire the map to the correct KQL database item id.
    kql_database_item_id = databases_item_ids[0]

    # Resolve actual KQL database name (displayName)
    get_db_url = f"https://api.fabric.microsoft.com/v1/workspaces/{cfg.workspace_id}/kqlDatabases/{kql_database_item_id}"
    db_resp = fabric.request("GET", get_db_url)
    db_resp.raise_for_status()

    kql_database_name = db_resp.json().get("displayName")
    if not kql_database_name:
        raise RuntimeError("KQL database response did not include displayName.")

    # Create the function that returns the latest location per vehicle.
    # Keep columns explicit so the map config can bind Latitude/Longitude.
    kql = f""".create-or-alter function with (skipvalidation=true) {cfg.kql_function_name}() {{
    table("{cfg.eventhouse_table_name}")
    | summarize arg_max(EventTime, *) by VehicleId
    | project Latitude, Longitude, VehicleId, EventTime
    }}"""

    mgmt_url = f"{query_service_uri}/v1/rest/mgmt"
    mgmt_payload = {"db": kql_database_name, "csl": kql}

    mgmt_resp = client.post(mgmt_url, headers=_kusto_headers(), json=mgmt_payload)

    if mgmt_resp.status_code >= 400:
        # Kusto usually returns a detailed JSON error body on 400s.
        raise RuntimeError(
            "Kusto mgmt call failed.\n"
            f"URL: {mgmt_url}\n"
            f"DB: {mgmt_payload.get('db')}\n"
            f"Status: {mgmt_resp.status_code}\n"
            f"Body: {mgmt_resp.text}"
        )

    print("KQL function created/updated:", cfg.kql_function_name)
    return kql_database_item_id

Note

Nazwy pól zwracane przez funkcję KQL muszą być zgodne z nazwami kolumn używanymi w definicji mapy (Latitude i Longitude w tym samouczku).

Kompilowanie map.json

build_map_json kompiluje i zwraca ładunek map.json definiujący zawartość mapy Fabric. Ładunek jest zgodny ze schematem definicji elementu mapy i składa się z czterech sekcji: dataSources (gdzie pochodzą dane), iconSources (opcjonalne znaczniki niestandardowe), layerSources (co jest wykonywane zapytanie i jak często) oraz layerSettings (sposób renderowania wyniku na mapie).

Na potrzeby tego samouczka dataSources wskazuje bazę danych KQL (itemType: "KqlDatabase") utworzoną wcześniej, a pojedynczy wpis w layerSources to warstwa oparta na Kusto (type: "kusto", queryType: "function"), której query wywołuje funkcję składowaną LatestVehicleLocations(). refreshIntervalMs jest odczytywany z cfg.refresh_interval_ms (domyślnie co 5000 ms), więc warstwa cyklicznie ponownie uruchamia funkcję, a mapa odzwierciedla napływ nowych danych niemal w czasie rzeczywistym.

Odpowiedni wpis layerSettings wiąże kolumny wynikowe warstwy z mapą za pomocą latitudeColumnName: "Latitude" i longitudeColumnName: "Longitude", przedstawia każdy wiersz jako punkt bubble oraz wyświetla VehicleId i EventTime w podpowiedziach. Funkcja wyświetla zmontowany ładunek, aby można było sprawdzić dokładny kod JSON, który wysyła wywołanie Create Map( Utwórz mapę).

Aby uzyskać więcej informacji na temat interfejsu API REST definicji mapy, zobacz Definicja elementu mapy.

Dodaj następujące elementy po create_kql_function funkcji:

# =========================================================
# Build map.json
# =========================================================

def build_map_json(cfg: Config, kql_database_item_id: str) -> dict:
    """
    Build and return the map.json payload for the Fabric Map.

    Wires `dataSources` to the KQL database created earlier, defines a
    single Kusto-backed layer in `layerSources` that calls the stored
    function `cfg.kql_function_name` and re-runs it every
    `cfg.refresh_interval_ms` milliseconds, and configures `layerSettings`
    to bind the `Latitude` / `Longitude` columns and render each row as a
    bubble point. Prints the assembled payload for inspection.
    """
    layer_source_id = str(uuid.uuid4())
    layer_setting_id = str(uuid.uuid4())
    data_source_name = "kqlConnection"

    map_json = {
        "$schema": "https://developer.microsoft.com/json-schemas/fabric/item/map/definition/2.0.0/schema.json",
        "basemap": {},

        "dataSources": [
            {
                "name": data_source_name,
                "itemType": "KqlDatabase",
                "workspaceId": cfg.workspace_id,
                "itemId": kql_database_item_id
            }
        ],

        "iconSources": [],

        "layerSources": [
            {
                "id": layer_source_id,
                "name": cfg.kql_function_name,
                "type": "kusto",
                "dataSourceName": data_source_name,
                "workspaceId": cfg.workspace_id,
                "itemId": kql_database_item_id,
                "refreshIntervalMs": cfg.refresh_interval_ms,
                "queryType": "function",
                "query": f"{cfg.kql_function_name}()"
            }
        ],
        "layerSettings": [
            {
                "id": layer_setting_id,
                "name": "Live Locations",
                "sourceId": layer_source_id,
                "options": {
                    "type": "vector",
                    "visible": True,
                    "pointLayerType": "bubble",
                    "tooltipKeys": ["VehicleId", "EventTime"],
                    "bubbleOptions": {
                        "color": "#0078D4"
                    }
                },
                "latitudeColumnName": "Latitude",
                "longitudeColumnName": "Longitude"

            }
        ]
    }

    
    print("Map definition (map.json):", json.dumps(map_json, indent=2))
    return map_json

Tworzenie platformy .platform (metadanych platformy)

build_platform_json kompiluje i zwraca opcjonalną .platform część, którą wywołanie Utwórz mapę może zawierać obok map.json , gdy chcesz ustawić metadane elementów innych niż domyślne na mapie. Dołączenie części .platform nie jest wymagane — Fabric stosuje domyślne metadane, gdy część zostanie pominięta — ale w tym samouczku pokazano, jak utworzyć element, aby można było ponownie użyć wzorca, gdy potrzebujesz jawnej kontroli nad typem elementu, nazwą wyświetlaną, opisem lub stabilnym identyfikatorem logicznym.

Ładunek jest zgodny ze schematem właściwości platformy i ma dwie sekcje: metadata (type: "Map", displayName, description) i config (version, logicalId). logicalId jest tutaj generowany jako nowy identyfikator UUID, co sprawdza się przy jednorazowym tworzeniu; jeśli planujesz ponownie wdrażać tę samą mapę przez integrację z Git lub podczas kolejnych uruchomień, ustaw logicalId na stałą wartość, aby aktualizacje były kierowane do tego samego elementu.

Aby uzyskać więcej informacji, zobacz Mapuj definicję elementu i Omówienie definicji elementu.

Dodaj następujące elementy po build_map_json funkcji:

# =========================================================
# Build .platform (platform metadata)
# =========================================================

def build_platform_json(cfg: Config) -> dict:
    """
    Build and return the optional .platform payload for a Fabric Map item.

    The map definition supports an optional .platform part alongside
    map.json that carries non-default item metadata: the item type,
    display name and description, and a `logicalId` used for
    deterministic updates. Fabric applies defaults when the part is
    omitted, so this payload is only needed when you want explicit
    control over those fields. A fresh UUID is used for `logicalId`
    here; pin it to a stable value if repeat runs should target the
    same item.
    """
    return {
        "$schema": "https://developer.microsoft.com/json-schemas/fabric/gitIntegration/platformProperties/2.0.0/schema.json",
        "metadata": {
            "type": "Map",
            "displayName": cfg.map_display_name,
            "description": cfg.map_description
        },
        "config": {
            "version": "2.0",
            # Use a stable logicalId if you want deterministic updates; UUID is fine for create.
            "logicalId": str(uuid.uuid4())
        }
    }

Utwórz mapę z definicją wbudowaną

create_map tworzy mapę, wysyłając metodą POST złożoną definicję inline, i zwraca identyfikator elementu nowej mapy. Żądanie zawiera trzy części zakodowane w formacie base64 w obszarze payloadType: "InlineBase64": map.json (wymagana definicja rdzenia), opcjonalne .platform metadane utworzone w poprzednim kroku oraz plik zapytania Kusto o nazwie queries/layerSource-<layerSourceId>.kql zawierający wywołanie przechowywanej funkcji KQL. Łączenie wszystkich trzech części w jednym wywołaniu aprowizuje mapę i podłącza warstwę danych do funkcji KQL niepodzielnie, więc nie jest wymagana kolejna getDefinition / updateDefinition runda.

Nazwa pliku zapytania ma znaczenie: Fabric odnajduje zapytanie warstwy, dopasowując queries/layerSource-<layerSourceId>.kql do id odpowiedniego wpisu w layerSources, więc funkcja pobiera identyfikator źródła warstwy z map_json["layerSources"][0]["id"], aby zbudować ścieżkę. map.json i .platform są zakodowane za pomocą _json_to_b64algorytmu base64; tekst zapytania jest zakodowany bezpośrednio w formacie base64, ponieważ jest to ciąg, a nie dict.

Interfejs API REST tworzenia mapy może odpowiedzieć za pomocą 201 Created (synchronicznego, wbudowanego identyfikatora), 202 Accepted (asynchronicznego LRO za pośrednictwem Location lub x-ms-operation-id), lub 200 OK z ładunkiem uzupełniania tylko do stanu, w którym mapa nie jest jeszcze widoczna w mapach listy z powodu opóźnienia propagacji zaplecza. _handle_lro obejmuje wszystkie te przypadki — w tym tworzenie listy i dopasowywanie według displayName — dlatego ta funkcja przekazuje jej pełną obsługę odpowiedzi w jednym wywołaniu.

Aby uzyskać więcej informacji, zobacz Definicja elementu mapy.

Dodaj następujące elementy po build_platform_json funkcji:

# =========================================================
# Create a map with inline definition
# =========================================================


def create_map(client: httpx.Client, fabric: FabricClient, cfg: Config, map_json: dict, platform_json: dict) -> str:
    """
    Create the Fabric Map with its definition inline and return its item ID.

    Sends a single Create Map request whose `parts` array carries three
    base64-encoded payloads: `map.json` (the required core definition),
    the optional `.platform` metadata, and a Kusto query file named
    `queries/layerSource-<layerSourceId>.kql` whose `<layerSourceId>`
    matches `map_json["layerSources"][0]["id"]` so Fabric can bind the
    query to the layer. Delegates response handling to `_handle_lro`,
    which covers synchronous, asynchronous, and status-only completions.
    """

    create_map_url = f"https://api.fabric.microsoft.com/v1/workspaces/{cfg.workspace_id}/maps"

    # Extract the layer source id so we can name the query file correctly
    layer_source_id = map_json["layerSources"][0]["id"]

    # Kusto query content (bind to the stored function)
    query_text = f"{cfg.kql_function_name}()"
    query_b64 = base64.b64encode(query_text.encode("utf-8")).decode("utf-8")

    create_map_payload = {
        "displayName": cfg.map_display_name,
        "description": cfg.map_description,
        "definition": {
            "parts": [
                {
                    "path": "map.json",
                    "payload": _json_to_b64(map_json),
                    "payloadType": "InlineBase64"
                },
                {
                    "path": ".platform",
                    "payload": _json_to_b64(platform_json),
                    "payloadType": "InlineBase64"
                },
                {
                    # Kusto layer query file naming convention
                    "path": f"queries/layerSource-{layer_source_id}.kql",
                    "payload": query_b64,
                    "payloadType": "InlineBase64"
                }
            ]
        }
    }

    map_resp = fabric.request("POST", create_map_url, json_body=create_map_payload)

    return _handle_lro(
        client, map_resp,
        list_url=create_map_url,
        match_display_name=cfg.map_display_name,
    )

Organizowanie przepływu pracy

main to jedyny punkt wejścia, który uruchamia samouczek od początku do końca. Tworzy instancję Config, otwiera jeden httpx.Client, współużywany przez wszystkie funkcje pomocnicze, opakowuje go w FabricClient, a następnie wywołuje każdą funkcję kroku w kolejności zależności: create_eventhousecreate_kql_table_if_missing (musi istnieć, zanim strumień zdarzeń powiąże się z nim) → create_eventstream_with_definitionseed_eventstream_from_csvverify_eventhouse_data (wychwytuje błędną konfigurację pozyskiwania, zanim rozpocznie się jakakolwiek praca związana z mapą) → create_kql_functionwait_for_kql_database_ready (mechanizm best-effort, aby funkcja Create Map mogła rozpoznać bazę danych KQL) → build_map_jsonbuild_platform_jsoncreate_map.

Kolejność ma znaczenie, ponieważ większość kroków wykorzystuje coś utworzonego w jednym z wcześniejszych kroków — create_eventstream_with_definition wymaga elementu databasesItemIds zasobu eventhouse, a create_map wymaga nazwy funkcji KQL oraz identyfikatora elementu bazy danych KQL. Ostatni blok print wyświetla identyfikatory wszystkich utworzonych zasobów, aby można było je znaleźć w portalu Fabric.

Dodaj następujące elementy po create_map funkcji:

# =========================================================
# main(): orchestrates the full workflow
# =========================================================

def main():
    """
    Orchestrate the tutorial workflow.

    1) Create eventhouse
    2) Create KQL table (required for ingestion)
    3) Create Eventstream (definition-based)
    4) Seed initial data so the map is not empty on first open
    5) Validate ingestion BEFORE moving on
    6) Create KQL function (required for Maps layer)
    7) Ensure KQL database is available to Maps
    8) Build map.json
    9) Build .platform metadata
    10) Create map with inline definition
    """
    cfg = Config()

    print("Initializing clients...")
    with httpx.Client(timeout=60) as client:
        fabric = FabricClient(client)

        # Step 1: Create eventhouse
        eventhouse_id = create_eventhouse(client, fabric, cfg)

        # Step 2: Ensure table exists BEFORE Eventstream binds to it
        create_kql_table_if_missing(client, fabric, cfg, eventhouse_id)

        # Step 3: Create Eventstream (definition-based)
        eventstream_id = create_eventstream_with_definition(client, fabric, cfg, eventhouse_id)

        # Step 4: Seed initial data so the map is not empty on first open
        seed_count = seed_eventstream_from_csv(cfg)

        # Step 5: Validate ingestion BEFORE moving on
        verify_eventhouse_data(client, fabric, cfg, eventhouse_id)

        # Step 6: Create KQL function (required for Maps layer)
        kql_database_item_id = create_kql_function(client, fabric, cfg, eventhouse_id)

        # Step 7: Ensure KQL database is available to Maps
        wait_for_kql_database_ready(client, cfg, kql_database_item_id)

        # Step 8: Build map.json (Kusto function layer)
        map_json = build_map_json(cfg, kql_database_item_id)

        # Step 9: Build .platform metadata
        platform_json = build_platform_json(cfg)

        # Step 10: Create map with inline definition
        map_id = create_map(client, fabric, cfg, map_json, platform_json)

        print("\nDONE")
        print("Eventhouse ID:", eventhouse_id)
        print("Eventstream ID:", eventstream_id)
        print(f"Seed events sent: {seed_count}")
        print("KQL database item ID:", kql_database_item_id)
        print("KQL function:", cfg.kql_function_name)
        print("Map ID:", map_id)

if __name__ == "__main__":
    main()

Uruchamianie aplikacji

Note

Eventhouse, eventstream, baza danych KQL i nazwy wyświetlane mapy muszą być unikatowe w obszarze roboczym. Przed ponownym uruchomieniem skryptu usuń elementy utworzone w poprzednim uruchomieniu z obszaru roboczego Fabric lub zmień odpowiednie nazwy wyświetlane w Config. W przeciwnym razie wywołania create kończą się błędem 409 ItemDisplayNameAlreadyInUse.

Podczas wykonywania skryptu zostanie wyświetlony monit z prośbą o wklejenie parametrów połączenia strumienia zdarzeń.

Aby pobrać tę wartość:

  1. Otwórz obszar roboczy Fabric
  2. Otwieranie strumienia zdarzeń utworzonego przez skrypt
  3. Wybierz Niestandardowe źródło punktu końcowego
  4. Otwórz uwierzytelnianie za pomocą klucza SAS
  5. Kopiuj parametry połączenia — klucz podstawowy

Zrzut ekranu przedstawiający obszar roboczy Microsoft Fabric z otwartym panelem uwierzytelniania za pomocą klucza SAS. Panel wyświetla pole Connection string-primary key, gotowe do skopiowania na potrzeby uwierzytelniania eventstream.

Wklej wartość do konsoli po wyświetleniu monitu.

Important

Skrypt wstrzymuje wykonywanie do momentu podania tej wartości.

Uruchom skrypt:

python create_realtime_map.py

Sprawdź, czy wszystkie elementy zostały utworzone:

Zrzut ekranu przedstawiający mapę ulicy Seattle z wieloma niebieskimi znacznikami lokalizacji skupionych w centrum i okolicznych dzielnicach. Panel warstw danych po lewej stronie zawiera włączoną warstwę Lokalizacji na żywo. Mapa wyświetla ulice, nazwy sąsiedztwa i przyciski sterowania po prawej stronie na potrzeby nawigacji i zarządzania warstwami.

Na tym etapie wszystkie zasoby są tworzone i konfigurowane.

Aby symulować ciągłe przesyłanie strumieniowe i obserwować aktualizację mapy w czasie zbliżonym do rzeczywistego, kontynuuj śledzenie Tutorial: symulowanie pozyskiwania danych w czasie rzeczywistym dla mapy przy użyciu interfejsów API REST i Python. Opiera się ona bezpośrednio na tym samouczku i ponownie wykorzystuje utworzone wcześniej obiekty: eventhouse, eventstream, funkcję KQL i mapę.

Podsumowanie

W tym samouczku udostępniono zasoby potrzebne do rozwiązania geoprzestrzennego działającego w czasie rzeczywistym w usłudze Microsoft Fabric przy użyciu interfejsów API REST platformy Fabric i języka Python.

Wykonano następujące czynności:

  • Utworzono bazę danych eventhouse i KQL przy użyciu interfejsu API REST Fabric
  • Utworzono strumień zdarzeń z niestandardowym punktem końcowym na potrzeby pozyskiwania zdarzeń przesyłania strumieniowego
  • Zdefiniowano funkcję KQL do wykonywania zapytań i kształtowania danych w czasie rzeczywistym na potrzeby wizualizacji mapy
  • Zbudowano i wdrożono mapę Fabric z osadzoną definicją odwołującą się do danych Eventhouse
  • Zainicjowano strumień zdarzeń początkowymi zdarzeniami, aby mapa od razu wyświetlała dane

Ta architektura przedstawia typowy wzorzec analizy w czasie rzeczywistym w Fabric:

  • Zewnętrzni producenci wysyłają zdarzenia do Eventstream
  • Eventstream kieruje i pozyskuje dane do Eventhouse
  • Funkcje KQL przekształcają dane
  • Mapy wysyłają zapytania do Eventhouse i odświeżają się automatycznie, aby odzwierciedlić nowe zdarzenia

Automatyzując tworzenie zasobów przy użyciu interfejsów API Python i REST, masz teraz powtarzalne podejście do tworzenia aplikacji przestrzennych w czasie rzeczywistym bez ręcznej konfiguracji. Aby przesyłać ciągłe dane do mapy, przejdź do kolejnego samouczka dotyczącego symulatora.

Następne kroki

Teraz, gdy rozumiesz pełny przepływ, możesz rozszerzyć to rozwiązanie w celu włączenia symulatora w czasie rzeczywistym.

Aby zapoznać się z samouczkiem przedstawiającym tworzenie symulatora w czasie rzeczywistym dla właśnie utworzonej mapy przy użyciu interfejsów API REST, zobacz: