Wysyłanie własnych trace'ów ETL i BI do Datadog APM

Datadog APM kojarzy się przede wszystkim z obserwowaniem aplikacji backendowych. Trace zaczyna się wraz z żądaniem HTTP, a kolejne spany pokazują wywołania usług, bazy danych i zewnętrznych API.
W projektach data sytuacja wygląda inaczej. Proces ETL może działać w hurtowni danych, usłudze zarządzanej albo zewnętrznym silniku przetwarzania, do którego nie można podłączyć tracera. Zwykle dostępna jest jedynie historia wykonań: identyfikator zadania, czas rozpoczęcia i zakończenia, status oraz komunikat błędu.
Takie rekordy można okresowo pobierać i wysyłać do Datadog jako własne trace’y i spany. Dzięki temu zespół data i BI zobaczy nie tylko liczbę błędów, ale również:
- całkowity czas od uruchomienia pipeline’u do odświeżenia raportu;
- czas poszczególnych etapów ekstrakcji, transformacji, ładowania i czekania;
- zapytania lub zadania, które spowalniają cały proces;
- zależności między pipeline’em, hurtownią i warstwą BI;
- miejsce oraz przyczynę przerwania przetwarzania.
Kiedy warto wysyłać procesy danych jako trace’y
To podejście ma sens, gdy:
- silnik ETL lub hurtownia udostępnia historię wykonań, ale nie pozwala na bezpośrednią instrumentację;
- rekord zawiera co najmniej czas rozpoczęcia, czas zakończenia i status zadania;
- ważna jest pełna ścieżka przetwarzania, a nie tylko pojedynczy pomiar;
- zespół chce analizować czasy etapów na flame graphie i szybko odnajdywać przyczynę opóźnienia raportu.
Jeżeli potrzebujesz wyłącznie wykresu czasu wykonania jednego joba, prostsza będzie metryka wysyłana przez DogStatsD. Trace’y są najbardziej przydatne wtedy, gdy jeden proces składa się z wielu powiązanych etapów (spanów).
Pobieranie historii wykonań
Kolektor może co minutę odpytywać tabelę systemową, API orkiestratora albo log zadań. Niezależnie od źródła dobrze jest sprowadzić rekordy do wspólnego modelu:
from dataclasses import dataclass
from datetime import datetime
@dataclass()
class JobExecution:
execution_id: str
pipeline_name: str
query: str
job_name: str
started_at: datetime
finished_at: datetime | None
status: str
error_message: str | None = None
rows_processed: int | None = None
Wysyłanie zakończonego joba jako span
ddtrace pozwala ustawić rzeczywisty czas rozpoczęcia i zakończenia spanu. Dzięki temu span może opisywać zadanie wykonane wcześniej w zewnętrznym systemie.
from ddtrace import tracer
def send_job_span(job: JobExecution) -> None:
if job.finished_at is None:
return
span = tracer.start_span(
"etl.job",
service="data-platform",
resource=job.job_name,
)
span.start = job.started_at.timestamp()
span.span_type = "worker"
span.set_tag("pipeline.name", job.pipeline_name)
span.set_tag("job.execution_id", job.execution_id)
span.set_tag("job.status", job.status)
if job.rows_processed is not None:
span.set_metric("rows_processed", job.rows_processed)
if job.status == "failed":
span.error = 1
span.set_tag("error.type", "SomeType")
span.set_tag(
"error.message",
job.error_message or f"Job {job.execution_id} failed",
)
span.finish(finish_time=job.finished_at.timestamp())
Warto rozdzielić znaczenie podstawowych pól:
serviceokreśla platformę lub logiczny system, na przykładdata-platform;- nazwa spanu opisuje rodzaj operacji, na przykład
etl.jobalbowarehouse.query; resourcewskazuje stabilną nazwę zadania, na przykładtransform_sales;- tagi przechowują identyfikator wykonania, nazwę pipeline’u, status i inne dane diagnostyczne;
- metryki spanu przechowują wartości liczbowe, na przykład liczbę przetworzonych wierszy.
Budowanie jednego trace’a dla całego pipeline’u
Największą wartość daje połączenie wszystkich etapów jednego przebiegu. Główny span opisuje pipeline, a spany poszczególnych jobów są jego dziećmi.
W rzeczywistym rozwiązaniu wszystkie rekordy trzeba grupować nie tylko po nazwie pipeline’u, lecz przede wszystkim po identyfikatorze jego konkretnego uruchomienia. Inaczej etapy z kilku równoległych przebiegów mogą trafić do jednego trace’a.
Powiązanie ETL z raportem BI
Jeśli narzędzie BI zapisuje identyfikator uruchomienia pipeline’u albo odziedziczony kontekst Datadog, odświeżenie zbioru danych można dołączyć do tego samego trace’a. Pozwala to odpowiedzieć na pytanie, czy raport był opóźniony przez transformację danych, ładowanie hurtowni czy samo odświeżenie modelu BI.
Kontekst trace’a może być przenoszony w metadanych zadania, tabeli audytowej, parametrze pipeline’u lub komentarzu SQL:
select /* TraceId: 4059091792851183189, SpanId: 13591579091350868334 */ ...
Po stronie kolektora identyfikatory można odtworzyć przez propagator Datadog:
from ddtrace.propagation.http import HTTPPropagator
parent_context = HTTPPropagator.extract(
{
"x-datadog-trace-id": str(trace_id),
"x-datadog-parent-id": str(span_id),
}
)
span = tracer.start_span(
"warehouse.query",
service="data-warehouse",
resource=query_name,
child_of=parent_context,
)
Jeżeli własny format integracji przechowuje wyłącznie 64-bitowe identyfikatory dziesiętne, obie strony muszą używać zgodnego formatu trace ID. Nie należy automatycznie wyłączać 128-bitowych identyfikatorów w całej organizacji — najpierw trzeba sprawdzić format propagowany przez używane wersje bibliotek i sposób zapisu kontekstu.
Ograniczenia rozwiązania
- Kolektor nie obserwuje pracy na żywo, lecz odtwarza ją na podstawie historii wykonań.
- Opóźnione odpytywanie oznacza opóźnione pojawienie się trace’a w Datadog.
- Agent i backend Datadog mogą odrzucić dane zbyt stare, dlatego trzeba monitorować opóźnienie kolektora.
- Każdy rekord zamieniony w span zwiększa wolumen danych APM. Nie każdy techniczny krok musi być osobnym spanem.
- Ponowne pobranie tego samego zakończonego rekordu może utworzyć duplikat. Kolektor powinien zapamiętywać ostatni przetworzony znacznik czasu lub identyfikatory wysłanych wykonań.
- Trace nie zastępuje tabeli audytowej. Datadog służy do obserwowalności i diagnostyki, a źródłowy system nadal pozostaje miejscem pełnej historii procesów.
Jeśli rozwiązanie ma być niezależne od jednego dostawcy, ten sam model można zbudować przy użyciu OpenTelemetry Python SDK. Datadog przyjmuje dane OTLP, dzięki czemu sposób opisywania pipeline’ów nie musi być trwale związany z biblioteką ddtrace.