const { BigQuery } = require('@google-cloud/bigquery') const { Storage } = require('@google-cloud/storage') const GoogleCloudAdapter = require('../GoogleCloudAdapter') const { ErrorUtils } = require('../../../utils/error') const $error = new ErrorUtils() const $storage = new Storage() // function mapOperatorToSql(operator: string) { // const operatorMap = { // equals: '=', // notEquals: '!=', // lessThan: '<', // lessOrEquals: '<=', // greaterThan: '>', // greaterOrEquals: '>=', // includes: 'LIKE', // notIncludes: 'NOT LIKE', // in: 'IN', // notIn: 'NOT IN', // // Add more operators as needed // } // if (operatorMap[operator as keyof typeof operatorMap] === undefined) { // throw new Error(`Operator ${operator} is not supported`) // } // return operatorMap[operator as keyof typeof operatorMap] // } // function formatValueForSql(value: any) { // // Add logic here to format the value based on its type (string, number, date, etc.) // if (typeof value === 'string') { // if (value.includes('(') && value.includes(')')) { // return value // } // return `'${value}'` // } // return value // } export class BigQueryAdapter extends GoogleCloudAdapter { constructor(props?: TBigQueryProps) { super() this.location = props?.location ?? 'us-central1' this.bigquery = new BigQuery(props) // this.parseFilterToSql = this.parseFilterToSql.bind(this) } async createDataset({ name }: TBigQueryCreateDatasetProps) { try { const [dataset] = await this.bigquery.dataset(name).create({ location: this.location, }) return dataset } catch (error) { throw $error.errorHandler({ error }) } } async getDataset({ name }: TBigQueryGetDatasetProps) { try { const [dataset] = await this.bigquery.dataset(name).get({ location: this.location, }) return dataset } catch (error) { throw $error.errorHandler({ error }) } } async deleteDataset({ name, force = true }: TBigQueryDeleteDatasetProps) { try { await this.bigquery .dataset(name) .delete({ location: this.location, force }) return true } catch (error) { throw $error.errorHandler({ error }) } } async loadFromBucket(props: TBigQueryLoadFromBucketProps) { try { const metadata = { sourceFormat: props.sourceFormat ?? 'CSV', skipLeadingRows: props.skipLeadingRows ?? 0, schema: { fields: props.schema, }, fieldDelimiter: '~', location: this.location, } const storageRef = $storage.bucket(props.bucketName).file(props.fileName) const [job] = await this.bigquery .dataset(props.dataset) .table(props.table) .load(storageRef, metadata) if (job.status.errors && job.status.errors.length > 0) { throw new Error(job.status.errors) } return job } catch (error) { throw $error.errorHandler({ error }) } } async runQuery(props: TBigQueryRunQueryProps) { try { const [job] = await this.bigquery.createQueryJob({ query: props.sql, location: this.location, }) const [rows] = await job.getQueryResults() return rows } catch (error) { throw $error.errorHandler({ error }) } } // parseFilterToSql(filter: any) { // const self = this // if (!filter) { // return '' // } // if (filter.or) { // const orConditions = filter.or.map(self.parseFilterToSql) // return `(${orConditions.join(' OR ')})` // } // if (filter.and) { // const andConditions = filter.and.map(self.parseFilterToSql) // return `(${andConditions.join(' AND ')})` // } // if (filter.operator && filter.field && filter.value) { // const { operator, field, value } = filter // return `${field} ${mapOperatorToSql(operator)} ${formatValueForSql(value)}` // } // return '' // } }