Skip to content

Commit 63322e7

Browse files
committed
feat: implement aggregation methods with support for grouping and median calculations in Clickhouse, MongoDB, and MySQL connectors
1 parent 5a00406 commit 63322e7

3 files changed

Lines changed: 243 additions & 4 deletions

File tree

‎adminforth/dataConnectors/clickhouse.ts‎

Lines changed: 68 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import { IAdminForthDataSourceConnector, IAdminForthSingleFilter, IAdminForthAndOrFilter, AdminForthResource, AdminForthResourceColumn } from '../types/Back.js';
1+
import { IAdminForthDataSourceConnector, IAdminForthSingleFilter, IAdminForthAndOrFilter, AdminForthResource, AdminForthResourceColumn, IAggregationRule, IGroupByRule, IGroupByDateTrunc, IGroupByField } from '../types/Back.js';
22
import AdminForthBaseConnector from './baseConnector.js';
33
import dayjs from 'dayjs';
44
import { createClient } from '@clickhouse/client'
@@ -444,13 +444,79 @@ class ClickhouseConnector extends AdminForthBaseConnector implements IAdminForth
444444
return { where, params };
445445
}
446446

447+
async getAggregateWithOriginalTypes({ resource, filters, aggregations, groupBy }: {
448+
resource: AdminForthResource;
449+
filters: IAdminForthAndOrFilter;
450+
aggregations: { [alias: string]: IAggregationRule };
451+
groupBy?: IGroupByRule;
452+
}): Promise<Array<{ group?: string; [key: string]: any }>> {
453+
454+
const tableName = `${this.dbName}.${resource.table}`;
455+
456+
const selectParts: string[] = [];
457+
let groupExpr: string | null = null;
458+
459+
if (groupBy?.type === 'date_trunc') {
460+
const g = groupBy as IGroupByDateTrunc;
461+
const tz = g.timezone ?? 'UTC';
462+
463+
const field = `toTimeZone(${g.field}, '${tz}')`;
464+
465+
switch (g.truncation) {
466+
case 'day': groupExpr = `toDate(toStartOfDay(${field}))`; break;
467+
case 'month': groupExpr = `toDate(toStartOfMonth(${field}))`; break;
468+
case 'week': groupExpr = `toDate(toStartOfWeek(${field}))`; break;
469+
case 'year': groupExpr = `toDate(toStartOfYear(${field}))`; break;
470+
}
471+
472+
selectParts.push(`${groupExpr} AS \`group\``);
473+
474+
} else if (groupBy?.type === 'field') {
475+
const g = groupBy as IGroupByField;
476+
groupExpr = `${g.field}`;
477+
selectParts.push(`${groupExpr} AS \`group\``);
478+
}
479+
480+
for (const [alias, rule] of Object.entries(aggregations)) {
481+
switch (rule.operation) {
482+
case 'count': selectParts.push(`count() AS \`${alias}\``); break;
483+
case 'sum': selectParts.push(`sum(${rule.field}) AS \`${alias}\``); break;
484+
case 'avg': selectParts.push(`avg(${rule.field}) AS \`${alias}\``); break;
485+
case 'min': selectParts.push(`min(${rule.field}) AS \`${alias}\``); break;
486+
case 'max': selectParts.push(`max(${rule.field}) AS \`${alias}\``); break;
487+
case 'median': selectParts.push(`quantile(0.5)(${rule.field}) AS \`${alias}\``); break;
488+
}
489+
}
490+
491+
const { where, params } = this.whereClause(resource, filters);
492+
493+
let query = `SELECT ${selectParts.join(', ')} FROM ${tableName} ${where}`;
494+
495+
if (groupExpr) {
496+
query += ` GROUP BY ${groupExpr} ORDER BY ${groupExpr} ASC`;
497+
}
498+
499+
const result = await this.client.query({
500+
query,
501+
format: 'JSONEachRow',
502+
query_params: params,
503+
});
504+
505+
const rows = await result.json();
506+
507+
return rows.map((r: any) => ({
508+
group: r.group,
509+
...r,
510+
}));
511+
}
512+
447513
async getDataWithOriginalTypes({ resource, limit, offset, sort, filters }: {
448514
resource: AdminForthResource,
449515
limit: number,
450516
offset: number,
451517
sort: { field: string, direction: AdminForthSortDirections }[],
452518
filters: IAdminForthAndOrFilter,
453-
}): Promise<any[]> {
519+
}): Promise<Array<{ group?: string, [key: string]: any }>> {
454520
const columns = resource.dataSourceColumns.map((col) => {
455521
// for decimal cast to string
456522
if (col.type == AdminForthDataTypes.DECIMAL) {

‎adminforth/dataConnectors/mongo.ts‎

Lines changed: 98 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
import dayjs from 'dayjs';
22
import { MongoClient } from 'mongodb';
33
import { Decimal128, Double } from 'bson';
4-
import { IAdminForthDataSourceConnector, IAdminForthSingleFilter, IAdminForthAndOrFilter, AdminForthResource } from '../types/Back.js';
4+
import { IAdminForthDataSourceConnector, IAdminForthSingleFilter, IAdminForthAndOrFilter, AdminForthResource, IAggregationRule, IGroupByRule, IGroupByDateTrunc, IGroupByField } from '../types/Back.js';
55
import AdminForthBaseConnector from './baseConnector.js';
66
import { afLogger } from '../modules/logger.js';
77
import { AdminForthDataTypes, AdminForthFilterOperators, AdminForthSortDirections, } from '../types/Common.js';
@@ -305,6 +305,103 @@ class MongoConnector extends AdminForthBaseConnector implements IAdminForthDataS
305305
.filter((f) => (f as IAdminForthSingleFilter).insecureRawSQL === undefined)
306306
.map((f) => this.getFilterQuery(resource, f)));
307307
}
308+
309+
async getAggregateWithOriginalTypes({ resource, filters, aggregations, groupBy }: {
310+
resource: AdminForthResource;
311+
filters: IAdminForthAndOrFilter;
312+
aggregations: any;
313+
groupBy?: any;
314+
}): Promise<Array<{ group?: string, [key: string]: any }>> {
315+
316+
const collection = this.client.db().collection(resource.table);
317+
318+
const match = filters?.subFilters?.length ? this.getFilterQuery(resource, filters) : {};
319+
320+
let groupId: any = null;
321+
322+
if (groupBy?.type === 'field') {
323+
groupId = `$${groupBy.field}`;
324+
}
325+
326+
if (groupBy?.type === 'date_trunc') {
327+
const tz = groupBy.timezone ?? 'UTC';
328+
329+
groupId = {
330+
$dateTrunc: {
331+
date: `$${groupBy.field}`,
332+
unit: groupBy.truncation,
333+
timezone: tz,
334+
},
335+
};
336+
}
337+
338+
const groupStage: any = {
339+
_id: groupId,
340+
};
341+
342+
for (const [alias, rule] of Object.entries(aggregations) as any) {
343+
switch (rule.operation) {
344+
case 'count': groupStage[alias] = { $sum: 1 }; break;
345+
case 'sum': groupStage[alias] = { $sum: { $toDouble: `$${rule.field}` } }; break;
346+
case 'avg': groupStage[alias] = { $avg: { $toDouble: `$${rule.field}` } }; break;
347+
case 'min': groupStage[alias] = { $min: { $toDouble: `$${rule.field}` } }; break;
348+
case 'max': groupStage[alias] = { $max: { $toDouble: `$${rule.field}` } }; break;
349+
case 'median': groupStage[alias] = { $push: { $toDouble: `$${rule.field}` } }; break;
350+
}
351+
}
352+
353+
const pipeline: any[] = [];
354+
355+
if (Object.keys(match).length) {
356+
pipeline.push({ $match: match });
357+
}
358+
359+
pipeline.push({ $group: groupStage });
360+
361+
pipeline.push({
362+
$project: {
363+
_id: 0,
364+
group: {
365+
$cond: {
366+
if: { $isNumber: "$_id" },
367+
then: "$_id",
368+
else: {
369+
$cond: {
370+
if: { $toString: "$_id" },
371+
then: { $dateToString: { format: "%Y-%m-%d", date: "$_id", timezone: groupBy?.timezone ?? 'UTC' } },
372+
else: "$_id"
373+
}
374+
}
375+
}
376+
},
377+
...Object.fromEntries(
378+
Object.keys(groupStage)
379+
.filter(k => k !== '_id')
380+
.map(k => [k, `$${k}`])
381+
),
382+
},
383+
});
384+
385+
const calculateMedian = (arr: any[]) => {
386+
if (!Array.isArray(arr) || arr.length === 0) return null;
387+
const sorted = [...arr].sort((a, b) => a - b);
388+
const mid = Math.floor(sorted.length / 2);
389+
return sorted.length % 2 === 0
390+
? (sorted[mid - 1] + sorted[mid]) / 2
391+
: sorted[mid];
392+
};
393+
394+
const result = await collection.aggregate(pipeline).toArray();
395+
396+
const medianAliases = Object.keys(aggregations).filter(alias => aggregations[alias].operation === 'median');
397+
398+
return result.map(row => {
399+
medianAliases.forEach(alias => {
400+
row[alias] = calculateMedian(row[alias]);
401+
});
402+
return row;
403+
});
404+
}
308405

309406
async getDataWithOriginalTypes({ resource, limit, offset, sort, filters }:
310407
{

‎adminforth/dataConnectors/mysql.ts‎

Lines changed: 77 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
import dayjs from 'dayjs';
2-
import { AdminForthResource, IAdminForthSingleFilter, IAdminForthAndOrFilter, IAdminForthDataSourceConnector, AdminForthConfig } from '../types/Back.js';
2+
import { AdminForthResource, IAdminForthSingleFilter, IAdminForthAndOrFilter, IAdminForthDataSourceConnector, AdminForthConfig, IAggregationRule, IGroupByRule, IGroupByDateTrunc, IGroupByField } from '../types/Back.js';
33
import { AdminForthDataTypes, AdminForthFilterOperators, AdminForthSortDirections, } from '../types/Common.js';
44
import AdminForthBaseConnector from './baseConnector.js';
55
import mysql from 'mysql2/promise';
@@ -338,6 +338,82 @@ class MysqlConnector extends AdminForthBaseConnector implements IAdminForthDataS
338338
} : { sql: '', values: [] };
339339
}
340340

341+
private calculateMedian(values: number[]): number | null {
342+
if (!values.length) return null;
343+
const sorted = values.sort((a, b) => a - b);
344+
const mid = Math.floor(sorted.length / 2);
345+
return sorted.length % 2 === 0
346+
? (sorted[mid - 1] + sorted[mid]) / 2
347+
: sorted[mid];
348+
}
349+
350+
async getAggregateWithOriginalTypes({ resource, filters, aggregations, groupBy }: {
351+
resource: AdminForthResource;
352+
filters: IAdminForthAndOrFilter;
353+
aggregations: { [alias: string]: IAggregationRule };
354+
groupBy?: IGroupByRule;
355+
}): Promise<Array<{ group?: string, [key: string]: any }>> {
356+
const tableName = resource.table;
357+
const selectParts: string[] = [];
358+
const medianAliases: string[] = [];
359+
let groupExpr: string | null = null;
360+
361+
if (groupBy?.type === 'field') {
362+
groupExpr = `\`${groupBy.field}\``;
363+
selectParts.push(`${groupExpr} AS \`group\``);
364+
} else if (groupBy?.type === 'date_trunc') {
365+
const g = groupBy as IGroupByDateTrunc;
366+
const tz = g.timezone ?? 'UTC';
367+
368+
const innerExpr = `COALESCE(CONVERT_TZ(\`${g.field}\`, 'UTC', '${tz}'), \`${g.field}\`)`;
369+
370+
switch (g.truncation) {
371+
case 'day': groupExpr = `DATE_FORMAT(${innerExpr}, '%Y-%m-%d')`; break;
372+
case 'month': groupExpr = `DATE_FORMAT(${innerExpr}, '%Y-%m-01')`; break;
373+
case 'year': groupExpr = `DATE_FORMAT(${innerExpr}, '%Y-01-01')`; break;
374+
case 'week': groupExpr = `DATE_FORMAT(DATE_SUB(${innerExpr}, INTERVAL WEEKDAY(${innerExpr}) DAY), '%Y-%m-%d')`; break;
375+
}
376+
377+
selectParts.push(`${groupExpr} AS \`group\``);
378+
}
379+
380+
for (const [alias, rule] of Object.entries(aggregations)) {
381+
const f = `\`${rule.field}\``;
382+
switch (rule.operation) {
383+
case 'sum': selectParts.push(`SUM(${f}) AS \`${alias}\``); break;
384+
case 'count': selectParts.push(`COUNT(*) AS \`${alias}\``); break;
385+
case 'avg': selectParts.push(`AVG(${f}) AS \`${alias}\``); break;
386+
case 'min': selectParts.push(`MIN(${f}) AS \`${alias}\``); break;
387+
case 'max': selectParts.push(`MAX(${f}) AS \`${alias}\``); break;
388+
case 'median':
389+
selectParts.push(`GROUP_CONCAT(${f}) AS \`${alias}\``);
390+
medianAliases.push(alias);
391+
break;
392+
}
393+
}
394+
395+
const { sql: where, values: filterValues } = this.whereClauseAndValues(filters);
396+
let query = `SELECT ${selectParts.join(', ')} FROM \`${tableName}\` ${where}`;
397+
if (groupExpr) query += ` GROUP BY ${groupExpr} ORDER BY ${groupExpr} ASC`;
398+
399+
await this.client.execute("SET SESSION group_concat_max_len = 1000000;");
400+
401+
const [rows]: any = await this.client.execute(query, filterValues);
402+
403+
if (medianAliases.length > 0) {
404+
return rows.map(row => {
405+
medianAliases.forEach(alias => {
406+
const raw = row[alias];
407+
const nums = raw ? raw.split(',').map(Number).filter(n => !isNaN(n)) : [];
408+
row[alias] = this.calculateMedian(nums);
409+
});
410+
return row;
411+
});
412+
}
413+
414+
return rows;
415+
}
416+
341417
async getDataWithOriginalTypes({ resource, limit, offset, sort, filters }): Promise<any[]> {
342418
const columns = resource.dataSourceColumns.map((col) => `${col.name}`).join(', ');
343419
const tableName = resource.table;

0 commit comments

Comments
 (0)