Monitorizarea proceselor ETL într-un mic depozit de date

Multe persoane folosesc instrumente specializate pentru a crea proceduri de extragere, transformare și încărcare a datelor în baze de date relaționale. Procesul de funcționare al instrumentelor este înregistrat, iar erorile sunt consemnate.

În cazul unei erori, jurnalul conține informații despre faptul că instrumentului nu i-a reușit să finalizeze sarcina și ce module (de obicei, java) s-au oprit. În ultimele linii se poate găsi eroarea bazei de date, de exemplu, o încălcare a cheii unice a tabelului.

Pentru a răspunde la întrebarea despre rolul informațiilor privind erorile ETL, am clasificat toate problemele apărute în ultimii doi ani într-un depozit semnificativ.

Monitorizarea proceselor ETL într-un mic depozit de date

Erorile bazei de date includ lucruri precum lipsa de spațiu, întreruperea conexiunii, sesiuni blocate și altele.

Erorile logice includ lucruri precum încălcarea cheilor tabelului, obiecte invalide, lipsa accesului la obiecte și altele.
Planificatorul poate fi pornit la timpul greșit, se poate bloca și altele.

Erorile simple nu necesită mult timp pentru a fi corectate. Majoritatea dintre ele sunt gestionate de un ETL bun în mod automat.

Erorile complexe necesită deschiderea și verificarea procedurilor de lucru cu datele, investigarea surselor de date. Acestea duc adesea la necesitatea testării modificărilor și desfășurării.

Așadar, jumătate din toate problemele sunt legate de baza de date. 48% din toate erorile sunt erori simple.
O treime dintre toate problemele sunt legate de schimbarea logicii sau a modelului depozitului, mai mult de jumătate dintre aceste erori fiind complexe.

Și mai puțin de un sfert din toate problemele sunt legate de planificatorul de sarcini, din care 18% sunt erori simple.

În general, 22% din toate erorile apărute sunt complexe, corectarea lor necesită cea mai mare atenție și timp. Acestea apar aproximativ o dată pe săptămână, în timp ce erorile simple se întâmplă aproape în fiecare zi.

Evident, monitorizarea proceselor ETL va fi eficientă atunci când în jurnal este indicat cât mai precis locul erorii și este necesar un timp minim pentru a localiza sursa problemei.

Monitorizare eficientă

Ce mi-ar plăcea să văd în procesul de monitorizare ETL?

Monitorizarea proceselor ETL într-un mic depozit de date
Start at — când a început lucrul,
Source — sursa datelor,
Layer — ce nivel al depozitului este încărcat,
ETL Job Name — procedura de încărcare care constă din numeroși pași mici,
Step Number — numărul pasului executat,
Rânduri afectate — câte date au fost deja procesate,
Durata sec — cât timp durează execuția,
Status — totul este bine sau nu: OK, ERROR, RUNNING, HANGS
Mesaj — ultimul mesaj reușit sau descrierea erorii.

Pe baza statusului înregistrărilor, se poate trimite un email altor participanți. Dacă nu sunt erori, atunci nici emailul nu este necesar.

Astfel, în caz de eroare, este clar indicat locul incidentului.

Uneori se întâmplă ca instrumentul de monitorizare să nu funcționeze. În acest caz, există posibilitatea de a apela direct în baza de date o vedere (view) pe baza căreia este construit raportul.

Tabelul de monitorizare ETL

Pentru a implementa monitorizarea proceselor ETL este suficient un singur tabel și o singură vedere.

Pentru aceasta, se poate întoarce la mica dvs. păstrătoare și crea un prototip în baza de date sqlite.

DDL tabelului

CREATE TABLE UTL_JOB_STATUS (

/* Tabel pentru logarea execuției muncii. Este important ca munca să aibă pașii ETL_START și ETL_END sau ETL_ERROR */
  UTL_JOB_STATUS_ID INTEGER NOT NULL PRIMARY KEY AUTOINCREMENT,
  SID               INTEGER NOT NULL DEFAULT -1, /* Identificator de sesiune. Unic pentru fiecare execuție a muncii */
  LOG_DT            INTEGER NOT NULL DEFAULT 0,  /* Dată și oră */
  LOG_D             INTEGER NOT NULL DEFAULT 0,  /* Dată */
  JOB_NAME          TEXT NOT NULL DEFAULT 'N/A', /* Numele muncii, de exemplu, JOB_STG2DM_GEO */
  STEP_NAME         TEXT NOT NULL DEFAULT 'N/A', /* ETL_START, ..., ETL_END/ETL_ERROR */
  STEP_DESCR        TEXT,                        /* Descrierea sarcinii sau mesajului de eroare */
  UNIQUE (SID, JOB_NAME, STEP_NAME)
);
INSERT INTO UTL_JOB_STATUS (UTL_JOB_STATUS_ID) VALUES (-1);

DDL vederii/raportului

CREAȚI VIZIUNE DACĂ NU EXISTĂ UTL_JOB_STATUS_V
CA 
/* Conținut: Jurnalul de execuție a pachetelor pentru ultimele 3 luni. */
CU SRC CA (
  SELECT LOG_D,
    LOG_DT,
    UTL_JOB_STATUS_ID,
    SID,
	CASE CÂND INSTR(JOB_NAME, 'FTP') AȘA 'TRANSFER' /* transfer de fișiere */
	     CÂND INSTR(JOB_NAME, 'STG') AȘA 'STADIA' /* stadiu */
	     CÂND INSTR(JOB_NAME, 'CLS') AȘA 'CURĂȚARE' /* curățare */
	     CÂND INSTR(JOB_NAME, 'DIM') AȘA 'DIMENSIUNE' /* dimensiune */
	     CÂND INSTR(JOB_NAME, 'FCT') AȘA 'FAPT' /* fapt */
		 CÂND INSTR(JOB_NAME, 'ETL') AȘA 'MART DE DATE' /* mart de date */
	     CÂND INSTR(JOB_NAME, 'RPT') AȘA 'REPORT' /* raport */
	     ALTCEVA 'N/A' FINE AS LAYER,
	CASE CÂND INSTR(JOB_NAME, 'ACCESS') AȘA 'JURNAL DE ACCES' /* sursă */
	     CÂND INSTR(JOB_NAME, 'MASTER') AȘA 'DATE PRINCIPALE' /* sursă */
	     CÂND INSTR(JOB_NAME, 'AD-HOC') AȘA 'AD-HOC' /* sursă */
	     ALTCEVA 'N/A' FINE AS SOURCE,
    JOB_NAME,
    STEP_NAME,
    CASE CÂND STEP_NAME='ETL_START' AȘA 1 ALTCEVA 0 FINE AS START_FLAG,
    CASE CÂND STEP_NAME='ETL_END' AȘA 1 ALTCEVA 0 FINE AS END_FLAG,
    CASE CÂND STEP_NAME='ETL_ERROR' AȘA 1 ALTCEVA 0 FINE AS ERROR_FLAG,
    STEP_NAME || ' : ' || STEP_DESCR AS STEP_LOG,
	SUBSTR( SUBSTR(STEP_DESCR, INSTR(STEP_DESCR, '***')+4), 1, INSTR(SUBSTR(STEP_DESCR, INSTR(STEP_DESCR, '***')+4), '***')-2 ) AS AFFECTED_ROWS
  FROM UTL_JOB_STATUS
  WHERE datetime(LOG_D, 'unixepoch') >= date('now', 'start of month', '-3 month')
)
SELECT JB.SID,
  JB.MIN_LOG_DT AS START_DT,
  strftime('%d.%m.%Y %H:%M', datetime(JB.MIN_LOG_DT, 'unixepoch')) AS LOG_DT,
  JB.SOURCE,
  JB.LAYER,
  JB.JOB_NAME,
  CASE
  CÂND JB.ERROR_FLAG = 1 AȘA 'EROARE'
  CÂND JB.ERROR_FLAG = 0 ȘI JB.END_FLAG = 0 ȘI strftime('%s','now') - JB.MIN_LOG_DT > 0.5*60*60 AȘA 'ÎN AȘTEPTARE' /* o jumătate de oră */
  CÂND JB.ERROR_FLAG = 0 ȘI JB.END_FLAG = 0 AȘA 'SE DESFĂȘOARĂ'
  ALTCEVA 'BINE'
  END AS STATUS,
  ERR.STEP_LOG     AS STEP_LOG,
  JB.CNT           AS STEP_CNT,
  JB.AFFECTED_ROWS AS AFFECTED_ROWS,
  strftime('%d.%m.%Y %H:%M', datetime(JB.MIN_LOG_DT, 'unixepoch')) AS JOB_START_DT,
  strftime('%d.%m.%Y %H:%M', datetime(JB.MAX_LOG_DT, 'unixepoch')) AS JOB_END_DT,
  JB.MAX_LOG_DT - JB.MIN_LOG_DT AS JOB_DURATION_SEC
FROM
  ( SELECT SID, SOURCE, LAYER, JOB_NAME,
           MAX(UTL_JOB_STATUS_ID) AS UTL_JOB_STATUS_ID,
           MAX(START_FLAG)       AS START_FLAG,
           MAX(END_FLAG)         AS END_FLAG,
           MAX(ERROR_FLAG)       AS ERROR_FLAG,
           MIN(LOG_DT)           AS MIN_LOG_DT,
           MAX(LOG_DT)           AS MAX_LOG_DT,
           SUM(1)                AS CNT,
           SUM(IFNULL(AFFECTED_ROWS, 0)) AS AFFECTED_ROWS
    FROM SRC
    GROUP BY SID, SOURCE, LAYER, JOB_NAME
  ) JB,
  ( SELECT UTL_JOB_STATUS_ID, SID, JOB_NAME, STEP_LOG
    FROM SRC
    WHERE 1 = 1
  ) ERR
WHERE 1 = 1
  ȘI JB.SID = ERR.SID
  ȘI JB.JOB_NAME = ERR.JOB_NAME
  ȘI JB.UTL_JOB_STATUS_ID = ERR.UTL_JOB_STATUS_ID
ORDER BY JB.MIN_LOG_DT DESC, JB.SID DESC, JB.SOURCE;

Verificare SQL pentru posibilitatea de a obține un nou număr de sesiune

SELECT SUM (
  CASE CÂND start_job.JOB_NAME NU ESTE NULL ȘI end_job.JOB_NAME ESTE NULL /* job existent s-a terminat */
	    ȘI NU ( 'y' = 'n' ) /* restart forțat PARAMETRU */
       AȘA 1 ALTCEVA 0
  END ) AS IS_RUNNING
  FROM
    ( SELECT 1 AS dummy FROM UTL_JOB_STATUS WHERE sid = -1) d_job
  LEFT OUTER JOIN
    ( SELECT JOB_NAME, SID, 1 AS dummy
      FROM UTL_JOB_STATUS
      WHERE JOB_NAME = 'RPT_ACCESS_LOG' /* numele jobului PARAMETRU */
	    ȘI STEP_NAME = 'ETL_START'
      GROUP BY JOB_NAME, SID
    ) start_job /* începe */
  PE d_job.dummy = start_job.dummy
  LEFT OUTER JOIN
    ( SELECT JOB_NAME, SID
      FROM UTL_JOB_STATUS
      WHERE JOB_NAME = 'RPT_ACCESS_LOG'  /* numele jobului PARAMETRU */
	    ȘI STEP_NAME in ('ETL_END', 'ETL_ERROR') /* statut de oprire */
      GROUP BY JOB_NAME, SID
    ) end_job /* se termină */
  PE start_job.JOB_NAME = end_job.JOB_NAME
     ȘI start_job.SID = end_job.SID

Caracteristici ale tabelului:

  • Începerea și finalizarea procedurii de prelucrare a datelor trebuie să fie însoțite de pașii ETL_START și ETL_END
  • În caz de eroare, trebuie să fie creat un pas ETL_ERROR cu descrierea acesteia
  • Numărul de date procesate trebuie evidențiat, de exemplu, cu ajutorul unor stele
  • În același timp, aceeași procedură poate fi lansată cu parametrul force_restart=y; fără acesta, numărul sesiunii este atribuit doar procedurii finalizate
  • În modul normal, nu se poate lansa simultan aceeași procedură de prelucrare a datelor

Operațiile necesare pentru lucrul cu tabelul sunt următoarele:

  • obținerea numărului sesiunii pentru procedura ETL lansată
  • inserarea unui înregistrări în jurnal în tabel
  • obținerea celei mai recente înregistrări de succes pentru procedura ETL

În baze de date precum Oracle sau Postgres, aceste operații pot fi implementate prin funcții încorporate. Pentru sqlite este necesar un mecanism extern, iar în acest caz, acesta a fost prototipat în PHP.

Ieșire

Astfel, mesajele de eroare în instrumentele de prelucrare a datelor joacă un rol mega-important. Totuși, pentru a găsi rapid cauza problemei, este dificil să le numim optime. Când numărul procedurilor se apropie de o sută, monitorizarea proceselor devine un proiect complex.

Articolul oferă un exemplu de soluție posibilă pentru problemă sub formă de prototip. Întregul prototip al unui mic depozit este disponibil pe gitlab Utilitare SQLite PHP ETL.

Sursa: habr.com

Cumpără un hosting fiabil pentru site-uri cu protecție DDoS, servere VPS VDS 🔥 Cumpără un hosting fiabil pentru site-uri cu protecție DDoS, servere VPS VDS | ProHoster