Skip to content

Commit bb9fe3b

Browse files
committed
test(data-quality): cover checks on real storage providers
1 parent 08daf0e commit bb9fe3b

8 files changed

Lines changed: 1004 additions & 0 deletions
Lines changed: 105 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,105 @@
1+
import { AthenaApiAdapter } from 'src/data-marts/data-storage-types/athena/adapters/athena-api.adapter';
2+
import { S3ApiAdapter } from 'src/data-marts/data-storage-types/athena/adapters/s3-api.adapter';
3+
import { AthenaConfig } from 'src/data-marts/data-storage-types/athena/schemas/athena-config.schema';
4+
import { AthenaCredentials } from 'src/data-marts/data-storage-types/athena/schemas/athena-credentials.schema';
5+
import { DataStorageType } from 'src/data-marts/data-storage-types/enums/data-storage-type.enum';
6+
import { registerRealDataQualitySuite } from './data-quality-real-suite';
7+
8+
const AWS_ACCESS_KEY_ID = process.env.AWS_ACCESS_KEY_ID;
9+
const AWS_SECRET_ACCESS_KEY = process.env.AWS_SECRET_ACCESS_KEY;
10+
const ATHENA_REGION = process.env.ATHENA_REGION;
11+
const ATHENA_OUTPUT_BUCKET = process.env.ATHENA_OUTPUT_BUCKET;
12+
const ATHENA_DATABASE = process.env.ATHENA_DATABASE;
13+
14+
const credentialsAvailable = Boolean(
15+
AWS_ACCESS_KEY_ID &&
16+
AWS_SECRET_ACCESS_KEY &&
17+
ATHENA_REGION &&
18+
ATHENA_OUTPUT_BUCKET &&
19+
ATHENA_DATABASE
20+
);
21+
22+
if (!credentialsAvailable) {
23+
console.log('Skipping Athena Data Quality integration tests: credentials are not configured');
24+
}
25+
26+
const describeIfCredentials = credentialsAvailable ? describe : describe.skip;
27+
const outputPrefix = `integration-test/data-quality-athena-${Date.now()}-${Math.random()
28+
.toString(36)
29+
.slice(2, 8)}/`;
30+
31+
describeIfCredentials('Athena Data Quality checks', () => {
32+
let adapter: AthenaApiAdapter;
33+
let s3Adapter: S3ApiAdapter;
34+
let config: AthenaConfig;
35+
36+
beforeAll(() => {
37+
const credentials: AthenaCredentials = {
38+
accessKeyId: AWS_ACCESS_KEY_ID!,
39+
secretAccessKey: AWS_SECRET_ACCESS_KEY!,
40+
};
41+
config = {
42+
region: ATHENA_REGION!,
43+
outputBucket: ATHENA_OUTPUT_BUCKET!,
44+
};
45+
adapter = new AthenaApiAdapter(credentials, config);
46+
s3Adapter = new S3ApiAdapter(credentials, config);
47+
});
48+
49+
afterAll(async () => {
50+
try {
51+
await s3Adapter.cleanupOutputFiles(config.outputBucket, outputPrefix);
52+
} catch (error) {
53+
console.warn('Failed to clean up Athena Data Quality query results:', error);
54+
}
55+
}, 90000);
56+
57+
registerRealDataQualitySuite({
58+
storageType: DataStorageType.AWS_ATHENA,
59+
schemaType: 'athena-data-mart-schema',
60+
nativeTypes: {
61+
integer: 'BIGINT',
62+
string: 'VARCHAR',
63+
timestamp: 'TIMESTAMP WITH TIME ZONE',
64+
},
65+
expressions: {
66+
integer: value => `CAST(${value ?? 'NULL'} AS BIGINT)`,
67+
string: value =>
68+
value === null
69+
? 'CAST(NULL AS VARCHAR)'
70+
: `CAST('${value.replaceAll("'", "''")}' AS VARCHAR)`,
71+
currentTimestamp: 'current_timestamp',
72+
staleTimestamp: "current_timestamp - INTERVAL '48' HOUR",
73+
},
74+
timeout: 180000,
75+
execute: async sql => {
76+
const { queryExecutionId } = await adapter.executeQuery(
77+
sql,
78+
config.outputBucket,
79+
outputPrefix
80+
);
81+
await adapter.waitForQueryToComplete(queryExecutionId);
82+
const result = await adapter.getQueryResults(queryExecutionId, undefined, 100);
83+
const columns = result.ResultSet?.ResultSetMetadata?.ColumnInfo ?? [];
84+
const rows = (result.ResultSet?.Rows ?? [])
85+
.slice(1)
86+
.map(row =>
87+
Object.fromEntries(
88+
columns.map((column, index) => [
89+
column.Name ?? `column_${index}`,
90+
row.Data?.[index]?.VarCharValue ?? null,
91+
])
92+
)
93+
);
94+
return {
95+
sql,
96+
rows,
97+
columnMetadata: columns.map(column => ({
98+
name: column.Name,
99+
label: column.Label,
100+
typeName: column.Type,
101+
})),
102+
};
103+
},
104+
});
105+
});
Lines changed: 74 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,74 @@
1+
import { BigQueryApiAdapter } from 'src/data-marts/data-storage-types/bigquery/adapters/bigquery-api.adapter';
2+
import { BigQueryServiceAccountCredentialsSchema } from 'src/data-marts/data-storage-types/bigquery/schemas/bigquery-credentials.schema';
3+
import {
4+
BIGQUERY_AUTODETECT_LOCATION,
5+
BigQueryConfig,
6+
} from 'src/data-marts/data-storage-types/bigquery/schemas/bigquery-config.schema';
7+
import { DataStorageType } from 'src/data-marts/data-storage-types/enums/data-storage-type.enum';
8+
import { registerRealDataQualitySuite } from './data-quality-real-suite';
9+
10+
const BQ_SERVICE_ACCOUNT_KEY = process.env.BQ_SERVICE_ACCOUNT_KEY;
11+
const BQ_PROJECT_ID = process.env.BQ_PROJECT_ID;
12+
const BQ_DATASET = process.env.BQ_DATASET;
13+
14+
const credentialsAvailable = Boolean(BQ_SERVICE_ACCOUNT_KEY && BQ_PROJECT_ID && BQ_DATASET);
15+
16+
export function registerBigQueryDataQualityIntegrationSuite(
17+
storageType: DataStorageType.GOOGLE_BIGQUERY | DataStorageType.LEGACY_GOOGLE_BIGQUERY,
18+
suiteName: string
19+
): void {
20+
if (!credentialsAvailable) {
21+
console.log(
22+
`Skipping ${suiteName} Data Quality integration tests: BigQuery credentials are not configured`
23+
);
24+
}
25+
26+
const describeIfCredentials = credentialsAvailable ? describe : describe.skip;
27+
describeIfCredentials(`${suiteName} Data Quality checks`, () => {
28+
let adapter: BigQueryApiAdapter;
29+
30+
beforeAll(() => {
31+
const credentials = BigQueryServiceAccountCredentialsSchema.parse(
32+
JSON.parse(BQ_SERVICE_ACCOUNT_KEY!)
33+
);
34+
const config: BigQueryConfig = {
35+
projectId: BQ_PROJECT_ID!,
36+
location: BIGQUERY_AUTODETECT_LOCATION,
37+
};
38+
adapter = new BigQueryApiAdapter(credentials, config);
39+
});
40+
41+
registerRealDataQualitySuite({
42+
storageType,
43+
schemaType: 'bigquery-data-mart-schema',
44+
nativeTypes: {
45+
integer: 'INT64',
46+
string: 'STRING',
47+
timestamp: 'TIMESTAMP',
48+
},
49+
expressions: {
50+
integer: value => `CAST(${value ?? 'NULL'} AS INT64)`,
51+
string: value =>
52+
value === null
53+
? 'CAST(NULL AS STRING)'
54+
: `CAST('${value.replaceAll("'", "''")}' AS STRING)`,
55+
currentTimestamp: 'CURRENT_TIMESTAMP()',
56+
staleTimestamp: 'TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 48 HOUR)',
57+
},
58+
fieldMode: 'NULLABLE',
59+
execute: async sql => {
60+
const { jobId } = await adapter.executeQuery(sql);
61+
const job = await adapter.getJob(jobId);
62+
const destination = job.metadata.configuration.query.destinationTable;
63+
if (!destination) throw new Error('BigQuery did not create a destination table');
64+
const table = adapter.createTableReference(
65+
destination.projectId,
66+
destination.datasetId,
67+
destination.tableId
68+
);
69+
const [rows] = await table.getRows({ maxResults: 100, autoPaginate: true });
70+
return { sql, rows: rows as Record<string, unknown>[] };
71+
},
72+
});
73+
});
74+
}
Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,4 @@
1+
import { DataStorageType } from 'src/data-marts/data-storage-types/enums/data-storage-type.enum';
2+
import { registerBigQueryDataQualityIntegrationSuite } from './data-quality-bigquery-test-support';
3+
4+
registerBigQueryDataQualityIntegrationSuite(DataStorageType.GOOGLE_BIGQUERY, 'BigQuery');
Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,70 @@
1+
import { DatabricksApiAdapter } from 'src/data-marts/data-storage-types/databricks/adapters/databricks-api.adapter';
2+
import { DatabricksAuthMethod } from 'src/data-marts/data-storage-types/databricks/enums/databricks-auth-method.enum';
3+
import { DatabricksConfig } from 'src/data-marts/data-storage-types/databricks/schemas/databricks-config.schema';
4+
import { DatabricksCredentials } from 'src/data-marts/data-storage-types/databricks/schemas/databricks-credentials.schema';
5+
import { DataStorageType } from 'src/data-marts/data-storage-types/enums/data-storage-type.enum';
6+
import { registerRealDataQualitySuite } from './data-quality-real-suite';
7+
8+
const DATABRICKS_HOST = process.env.DATABRICKS_HOST;
9+
const DATABRICKS_HTTP_PATH = process.env.DATABRICKS_HTTP_PATH;
10+
const DATABRICKS_TOKEN = process.env.DATABRICKS_TOKEN;
11+
const DATABRICKS_CATALOG = process.env.DATABRICKS_CATALOG;
12+
const DATABRICKS_SCHEMA = process.env.DATABRICKS_SCHEMA;
13+
14+
const credentialsAvailable = Boolean(
15+
DATABRICKS_HOST &&
16+
DATABRICKS_HTTP_PATH &&
17+
DATABRICKS_TOKEN &&
18+
DATABRICKS_CATALOG &&
19+
DATABRICKS_SCHEMA
20+
);
21+
22+
if (!credentialsAvailable) {
23+
console.log('Skipping Databricks Data Quality integration tests: credentials are not configured');
24+
}
25+
26+
const describeIfCredentials = credentialsAvailable ? describe : describe.skip;
27+
28+
describeIfCredentials('Databricks Data Quality checks', () => {
29+
let adapter: DatabricksApiAdapter;
30+
31+
beforeAll(() => {
32+
const credentials: DatabricksCredentials = {
33+
authMethod: DatabricksAuthMethod.PERSONAL_ACCESS_TOKEN,
34+
token: DATABRICKS_TOKEN!,
35+
};
36+
const config: DatabricksConfig = {
37+
host: DATABRICKS_HOST!,
38+
httpPath: DATABRICKS_HTTP_PATH!,
39+
};
40+
adapter = new DatabricksApiAdapter(credentials, config);
41+
});
42+
43+
afterAll(async () => {
44+
await adapter.destroy();
45+
}, 60000);
46+
47+
registerRealDataQualitySuite({
48+
storageType: DataStorageType.DATABRICKS,
49+
schemaType: 'databricks-data-mart-schema',
50+
nativeTypes: {
51+
integer: 'BIGINT',
52+
string: 'STRING',
53+
timestamp: 'TIMESTAMP',
54+
},
55+
expressions: {
56+
integer: value => `CAST(${value ?? 'NULL'} AS BIGINT)`,
57+
string: value =>
58+
value === null
59+
? 'CAST(NULL AS STRING)'
60+
: `CAST('${value.replaceAll("'", "''")}' AS STRING)`,
61+
currentTimestamp: 'current_timestamp()',
62+
staleTimestamp: 'current_timestamp() - INTERVAL 48 HOURS',
63+
},
64+
timeout: 120000,
65+
execute: async sql => ({
66+
sql,
67+
rows: await adapter.executeQueryAndFetchAll(sql),
68+
}),
69+
});
70+
});
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
import { DataStorageType } from 'src/data-marts/data-storage-types/enums/data-storage-type.enum';
2+
import { registerBigQueryDataQualityIntegrationSuite } from './data-quality-bigquery-test-support';
3+
4+
registerBigQueryDataQualityIntegrationSuite(
5+
DataStorageType.LEGACY_GOOGLE_BIGQUERY,
6+
'Legacy BigQuery'
7+
);

0 commit comments

Comments
 (0)