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.

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?

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 ș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.SIDCaracteristici 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 .
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 .
Sursa: habr.com
