gl aggregate table (ETL), cronjob by NODEJS
This commit is contained in:
@@ -0,0 +1,172 @@
|
||||
<?php
|
||||
/**
|
||||
* EtlManager
|
||||
*
|
||||
* Manages etl_gl_summary — the materialized aggregation of td_gl_item by period.
|
||||
* Used for fast dashboard queries instead of scanning raw GL rows.
|
||||
*
|
||||
* rebuildFull() — wipe and rebuild all periods for this company (manual trigger)
|
||||
* rebuildPeriod() — rebuild a single period (nightly scheduler, gap repair)
|
||||
* getStatus() — return last run info from company_setting
|
||||
*/
|
||||
class EtlManager
|
||||
{
|
||||
private PDO $pdo;
|
||||
private PDO $pdo1;
|
||||
private int $companyId;
|
||||
|
||||
public function __construct(PDO $pdo, PDO $pdo1, int $company_id)
|
||||
{
|
||||
$this->pdo = $pdo;
|
||||
$this->pdo1 = $pdo1;
|
||||
$this->companyId = $company_id;
|
||||
}
|
||||
|
||||
// ── Public ───────────────────────────────────────────────────────────────
|
||||
|
||||
public function checkAndRepair(): array
|
||||
{
|
||||
$gaps = $this->getGaps();
|
||||
|
||||
foreach ($gaps as $period) {
|
||||
$this->rebuildPeriod($period);
|
||||
}
|
||||
|
||||
$gaps_found = count($gaps);
|
||||
$periods_rebuilt = $gaps_found;
|
||||
$this->saveStatus('ok', $gaps_found, $periods_rebuilt);
|
||||
|
||||
return ['gaps_found' => $gaps_found, 'periods_rebuilt' => $periods_rebuilt];
|
||||
}
|
||||
|
||||
public function getGaps(): array
|
||||
{
|
||||
$sth = $this->pdo->prepare("
|
||||
SELECT DISTINCT g.period
|
||||
FROM td_gl g
|
||||
LEFT JOIN (
|
||||
SELECT period, MIN(source_updated_at) AS etl_ts
|
||||
FROM etl_gl_summary
|
||||
WHERE company_id = :cid2
|
||||
GROUP BY period
|
||||
) etl ON etl.period = g.period
|
||||
WHERE g.company_id = :cid
|
||||
AND g.source_type NOT IN ('voided', 'reversal')
|
||||
AND (
|
||||
etl.period IS NULL
|
||||
OR g.updated_at > etl.etl_ts
|
||||
)
|
||||
ORDER BY g.period
|
||||
");
|
||||
$sth->execute([':cid' => $this->companyId, ':cid2' => $this->companyId]);
|
||||
return $sth->fetchAll(PDO::FETCH_COLUMN);
|
||||
}
|
||||
|
||||
public function rebuildFull(): array
|
||||
{
|
||||
$this->pdo->prepare(
|
||||
"DELETE FROM etl_gl_summary WHERE company_id = :cid"
|
||||
)->execute([':cid' => $this->companyId]);
|
||||
|
||||
$this->pdo->prepare("
|
||||
INSERT INTO etl_gl_summary
|
||||
(company_id, acc_code, period, debit, credit, count, source_updated_at, updated_at)
|
||||
SELECT
|
||||
i.company_id,
|
||||
i.account_code,
|
||||
g.period,
|
||||
SUM(i.debit),
|
||||
SUM(i.credit),
|
||||
COUNT(*),
|
||||
MAX(g.updated_at),
|
||||
NOW()
|
||||
FROM td_gl_item i
|
||||
JOIN td_gl g ON g.id = i.gl_id AND g.company_id = i.company_id
|
||||
WHERE i.company_id = :cid
|
||||
AND g.source_type != 'voided'
|
||||
GROUP BY i.company_id, i.account_code, g.period
|
||||
")->execute([':cid' => $this->companyId]);
|
||||
|
||||
$periods_rebuilt = (int)$this->pdo->query(
|
||||
"SELECT COUNT(DISTINCT period) FROM etl_gl_summary WHERE company_id = {$this->companyId}"
|
||||
)->fetchColumn();
|
||||
|
||||
$this->saveStatus('ok', 0, $periods_rebuilt);
|
||||
|
||||
return ['periods_rebuilt' => $periods_rebuilt];
|
||||
}
|
||||
|
||||
public function rebuildPeriod(string $period): void
|
||||
{
|
||||
$this->pdo->prepare(
|
||||
"DELETE FROM etl_gl_summary WHERE company_id = :cid AND period = :period"
|
||||
)->execute([':cid' => $this->companyId, ':period' => $period]);
|
||||
|
||||
$this->pdo->prepare("
|
||||
INSERT INTO etl_gl_summary
|
||||
(company_id, acc_code, period, debit, credit, count, source_updated_at, updated_at)
|
||||
SELECT
|
||||
i.company_id,
|
||||
i.account_code,
|
||||
g.period,
|
||||
SUM(i.debit),
|
||||
SUM(i.credit),
|
||||
COUNT(*),
|
||||
MAX(g.updated_at),
|
||||
NOW()
|
||||
FROM td_gl_item i
|
||||
JOIN td_gl g ON g.id = i.gl_id AND g.company_id = i.company_id
|
||||
WHERE i.company_id = :cid
|
||||
AND g.period = :period
|
||||
AND g.source_type != 'voided'
|
||||
GROUP BY i.company_id, i.account_code, g.period
|
||||
")->execute([':cid' => $this->companyId, ':period' => $period]);
|
||||
}
|
||||
|
||||
public function getStatus(): array
|
||||
{
|
||||
$keys = ['etl_last_ran', 'etl_last_result', 'etl_gaps_found', 'etl_periods_rebuilt'];
|
||||
$sth = $this->pdo1->prepare(
|
||||
"SELECT setting_key, value FROM company_setting
|
||||
WHERE company_id = :cid AND setting_key IN ('" . implode("','", $keys) . "')"
|
||||
);
|
||||
$sth->execute([':cid' => $this->companyId]);
|
||||
$rows = $sth->fetchAll(PDO::FETCH_KEY_PAIR);
|
||||
|
||||
$slot_hour = $this->companyId % 24;
|
||||
$last_ran = $rows['etl_last_ran'] ?? null;
|
||||
$next_run = $last_ran
|
||||
? date('Y-m-d', strtotime($last_ran . ' +1 day')) . sprintf(' %02d:00', $slot_hour)
|
||||
: date('Y-m-d') . sprintf(' %02d:00', $slot_hour);
|
||||
|
||||
return [
|
||||
'slot_hour' => $slot_hour,
|
||||
'last_ran' => $last_ran,
|
||||
'last_result' => $rows['etl_last_result'] ?? null,
|
||||
'gaps_found' => (int)($rows['etl_gaps_found'] ?? 0),
|
||||
'periods_rebuilt' => (int)($rows['etl_periods_rebuilt'] ?? 0),
|
||||
'next_run' => $next_run,
|
||||
];
|
||||
}
|
||||
|
||||
// ── Private ──────────────────────────────────────────────────────────────
|
||||
|
||||
private function saveStatus(string $result, int $gaps_found, int $periods_rebuilt): void
|
||||
{
|
||||
$now = date('Y-m-d H:i:s');
|
||||
$data = [
|
||||
'etl_last_ran' => $now,
|
||||
'etl_last_result' => $result,
|
||||
'etl_gaps_found' => (string)$gaps_found,
|
||||
'etl_periods_rebuilt' => (string)$periods_rebuilt,
|
||||
];
|
||||
$sth = $this->pdo1->prepare(
|
||||
"INSERT INTO company_setting (company_id, setting_key, value, updated_at)
|
||||
VALUES (:cid, :key, :val, :ts)
|
||||
ON DUPLICATE KEY UPDATE value = VALUES(value), updated_at = VALUES(updated_at)"
|
||||
);
|
||||
foreach ($data as $key => $val) {
|
||||
$sth->execute([':cid' => $this->companyId, ':key' => $key, ':val' => $val, ':ts' => $now]);
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user