mirror of
https://github.com/fluxerapp/fluxer
synced 2026-10-11 13:12:23 +09:00
perf(api): stop Postgres KV reads scanning whole tables (#2118)
This commit is contained in:
1 parent
990176ac7c
commit
14ae64f5f3
11 files changed
+5348
-65
No files matched your search
@@ -2,11 +2,14 @@
|
||||
|
||||
import {type IPostgresClient, type PostgresQueryable, quoteIdentifier} from '@pkgs/postgres/src/Client';
|
||||
import cassandra from 'cassandra-driver';
|
||||
import {Logger} from '../Logger';
|
||||
import {getKvMeta, getTableMetadata} from './CassandraMetaRegistry';
|
||||
import type {CassandraParams, ColumnName, KvQueryMeta, PreparedQuery, WhereExpr} from './CassandraTypes';
|
||||
|
||||
type Row = Record<string, unknown>;
|
||||
type EqWhereExpr = Extract<WhereExpr<Row>, {kind: 'eq'}>;
|
||||
type InWhereExpr = Extract<WhereExpr<Row>, {kind: 'in'}>;
|
||||
type PinnedWhereExpr = EqWhereExpr | InWhereExpr;
|
||||
|
||||
interface StoredRow {
|
||||
row_key: string;
|
||||
@@ -17,7 +20,40 @@ interface PageState {
|
||||
offset: number;
|
||||
}
|
||||
|
||||
export type CandidatePlan =
|
||||
| {kind: 'none'}
|
||||
| {kind: 'rowKeys'; rowKeys: Array<string>}
|
||||
| {kind: 'range'; lowerBound: string; upperBound: string}
|
||||
| {kind: 'ranges'; lowerBounds: Array<string>; upperBounds: Array<string>}
|
||||
| {kind: 'partitionKeys'; partitionKeys: Array<string>}
|
||||
| {kind: 'scan'};
|
||||
|
||||
export interface QueryPlan {
|
||||
candidates: CandidatePlan;
|
||||
exact: boolean;
|
||||
}
|
||||
|
||||
interface QueryShape {
|
||||
leadingClauses: ReadonlyArray<PinnedWhereExpr>;
|
||||
partitionClauses: ReadonlyArray<PinnedWhereExpr> | null;
|
||||
inClauses: ReadonlyArray<InWhereExpr>;
|
||||
whereColumns: ReadonlyArray<string>;
|
||||
requiredColumns: ReadonlyArray<string> | null;
|
||||
clauseCount: number;
|
||||
summary: string;
|
||||
}
|
||||
|
||||
export interface PlanFragments {
|
||||
predicate: string;
|
||||
params: Array<unknown>;
|
||||
}
|
||||
|
||||
const VALUE_SEPARATOR = '\u001f';
|
||||
const KEY_RANGE_UPPER = ' ';
|
||||
const MAX_ROW_KEY_COMBINATIONS = 32_768;
|
||||
const MAX_PREFIX_RANGES = 256;
|
||||
const FULL_SCAN_LOG_INTERVAL_MS = 60_000;
|
||||
const FULL_SCAN_LOG_KEY_LIMIT = 1024;
|
||||
const ENCODED_TYPE_KEY = '__fluxer_type';
|
||||
const POSTGRES_KV_SCHEMA_LOCK_NAMESPACE = 0x46584b56;
|
||||
const POSTGRES_KV_SCHEMA_LOCK_TIMEOUT = '120s';
|
||||
@@ -93,11 +129,24 @@ function decodeRow(value: unknown): Row {
|
||||
return decoded;
|
||||
}
|
||||
|
||||
function decodeRowColumns(value: unknown, columns: ReadonlyArray<string>): Row {
|
||||
if (!isPlainObject(value)) {
|
||||
throw new Error('Postgres KV row payload is not an object');
|
||||
}
|
||||
const decoded: Row = {};
|
||||
for (const column of columns) {
|
||||
if (column in value) {
|
||||
decoded[column] = decodeValue(value[column]);
|
||||
}
|
||||
}
|
||||
return decoded;
|
||||
}
|
||||
|
||||
function valueKey(value: unknown): string {
|
||||
return JSON.stringify(encodeValue(value));
|
||||
}
|
||||
|
||||
function keyFromColumns(columns: ReadonlyArray<string>, row: Row): string {
|
||||
export function keyFromColumns(columns: ReadonlyArray<string>, row: Row): string {
|
||||
return columns.map((column) => valueKey(row[column])).join(VALUE_SEPARATOR);
|
||||
}
|
||||
|
||||
@@ -137,10 +186,6 @@ function rowKeyFromParams(meta: KvQueryMeta, params: CassandraParams): string {
|
||||
return rowKey(meta, paramsRow(params, (meta.pkColumns ?? meta.table.primaryKey) as ReadonlyArray<string>));
|
||||
}
|
||||
|
||||
function partitionKeyFromParams(meta: KvQueryMeta, params: CassandraParams): string {
|
||||
return partitionKey(meta, paramsRow(params, meta.table.partitionKey as ReadonlyArray<string>));
|
||||
}
|
||||
|
||||
function compareValues(left: unknown, right: unknown): number {
|
||||
if (typeof left === 'bigint' || typeof right === 'bigint') {
|
||||
const l = typeof left === 'bigint' ? left : BigInt(left as number | string);
|
||||
@@ -160,8 +205,8 @@ function valuesEqual(left: unknown, right: unknown): boolean {
|
||||
if (left == null && right == null) return true;
|
||||
if (left instanceof Date && right instanceof Date) return left.getTime() === right.getTime();
|
||||
if (Buffer.isBuffer(left) && Buffer.isBuffer(right)) return left.equals(right);
|
||||
if (left?.constructor?.name === 'LocalDate' || right?.constructor?.name === 'LocalDate') {
|
||||
return left?.toString() === right?.toString();
|
||||
if (left?.constructor?.name === 'LocalDate' && right?.constructor?.name === 'LocalDate') {
|
||||
return left.toString() === right.toString();
|
||||
}
|
||||
return left === right;
|
||||
}
|
||||
@@ -170,7 +215,11 @@ function getParam(params: CassandraParams, param: string): unknown {
|
||||
return params[param];
|
||||
}
|
||||
|
||||
function matchesWhere(row: Row, where: ReadonlyArray<WhereExpr<Row>> | undefined, params: CassandraParams): boolean {
|
||||
export function matchesWhere(
|
||||
row: Row,
|
||||
where: ReadonlyArray<WhereExpr<Row>> | undefined,
|
||||
params: CassandraParams,
|
||||
): boolean {
|
||||
for (const clause of where ?? []) {
|
||||
switch (clause.kind) {
|
||||
case 'eq':
|
||||
@@ -238,39 +287,269 @@ function sortRows(meta: KvQueryMeta, rows: Array<Row>): Array<Row> {
|
||||
});
|
||||
}
|
||||
|
||||
function equalityParam(where: ReadonlyArray<WhereExpr<Row>> | undefined, column: string): string | null {
|
||||
const clause = (where ?? []).find((entry) => entry.kind === 'eq' && entry.col === column);
|
||||
return clause && clause.kind === 'eq' ? clause.param : null;
|
||||
function whereClauses(meta: KvQueryMeta): ReadonlyArray<WhereExpr<Row>> {
|
||||
return (meta.where ?? []) as ReadonlyArray<WhereExpr<Row>>;
|
||||
}
|
||||
|
||||
function inParam(where: ReadonlyArray<WhereExpr<Row>> | undefined, column: string): string | null {
|
||||
const clause = (where ?? []).find((entry) => entry.kind === 'in' && entry.col === column);
|
||||
return clause && clause.kind === 'in' ? clause.param : null;
|
||||
}
|
||||
|
||||
function fullRowKeysFromWhere(meta: KvQueryMeta, params: CassandraParams): Array<string> | null {
|
||||
const pk = meta.table.primaryKey as ReadonlyArray<string>;
|
||||
const eqParams = pk.map((column) => equalityParam(meta.where as ReadonlyArray<WhereExpr<Row>> | undefined, column));
|
||||
if (eqParams.every((param) => param !== null)) {
|
||||
const row: Row = {};
|
||||
for (let i = 0; i < pk.length; i += 1) row[pk[i]!] = params[eqParams[i]!];
|
||||
return [rowKey(meta, row)];
|
||||
}
|
||||
if (pk.length === 1) {
|
||||
const param = inParam(meta.where as ReadonlyArray<WhereExpr<Row>> | undefined, pk[0]!);
|
||||
if (param) {
|
||||
const values = params[param] as ReadonlyArray<unknown> | Set<unknown>;
|
||||
const haystack = values instanceof Set ? [...values] : values;
|
||||
return haystack.map((value) => rowKey(meta, {[pk[0]!]: value}));
|
||||
}
|
||||
function pinnedClause(where: ReadonlyArray<WhereExpr<Row>>, column: string): PinnedWhereExpr | null {
|
||||
for (const clause of where) {
|
||||
if ((clause.kind === 'eq' || clause.kind === 'in') && clause.col === column) return clause;
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
function hasFullPartition(meta: KvQueryMeta): boolean {
|
||||
return meta.table.partitionKey.every((column) =>
|
||||
equalityParam(meta.where as ReadonlyArray<WhereExpr<Row>> | undefined, column),
|
||||
);
|
||||
function clauseColumns(clause: WhereExpr<Row>): ReadonlyArray<string> {
|
||||
return clause.kind === 'tupleGt' ? (clause.cols as ReadonlyArray<string>) : [clause.col as string];
|
||||
}
|
||||
|
||||
function describeClause(clause: WhereExpr<Row>): string {
|
||||
return `${clauseColumns(clause).join('+')} ${clause.kind}`;
|
||||
}
|
||||
|
||||
function requiredColumnsFor(meta: KvQueryMeta, whereColumns: ReadonlyArray<string>): ReadonlyArray<string> | null {
|
||||
if (!meta.columns) return null;
|
||||
const required = new Set<string>(meta.columns as ReadonlyArray<string>);
|
||||
for (const column of whereColumns) required.add(column);
|
||||
if (meta.orderBy) required.add(meta.orderBy.col as string);
|
||||
for (const column of meta.table.primaryKey as ReadonlyArray<string>) required.add(column);
|
||||
for (const column of meta.table.columns as ReadonlyArray<string>) {
|
||||
if (!required.has(column)) return [...required];
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
const QUERY_SHAPES = new WeakMap<KvQueryMeta, QueryShape>();
|
||||
|
||||
function queryShape(meta: KvQueryMeta): QueryShape {
|
||||
const cached = QUERY_SHAPES.get(meta);
|
||||
if (cached) return cached;
|
||||
const where = whereClauses(meta);
|
||||
const leadingClauses: Array<PinnedWhereExpr> = [];
|
||||
for (const column of meta.table.primaryKey as ReadonlyArray<string>) {
|
||||
const clause = pinnedClause(where, column);
|
||||
if (!clause) break;
|
||||
leadingClauses.push(clause);
|
||||
}
|
||||
const partitionColumns = meta.table.partitionKey as ReadonlyArray<string>;
|
||||
const partitionClauses: Array<PinnedWhereExpr> = [];
|
||||
for (const column of partitionColumns) {
|
||||
const clause = pinnedClause(where, column);
|
||||
if (!clause) {
|
||||
partitionClauses.length = 0;
|
||||
break;
|
||||
}
|
||||
partitionClauses.push(clause);
|
||||
}
|
||||
const whereColumns = [...new Set(where.flatMap(clauseColumns))];
|
||||
const shape: QueryShape = {
|
||||
leadingClauses,
|
||||
partitionClauses:
|
||||
partitionColumns.length > 0 && partitionClauses.length === partitionColumns.length ? partitionClauses : null,
|
||||
inClauses: where.filter((clause): clause is InWhereExpr => clause.kind === 'in'),
|
||||
whereColumns,
|
||||
requiredColumns: requiredColumnsFor(meta, whereColumns),
|
||||
clauseCount: where.length,
|
||||
summary: where.map(describeClause).join(', '),
|
||||
};
|
||||
QUERY_SHAPES.set(meta, shape);
|
||||
return shape;
|
||||
}
|
||||
|
||||
function clauseValues(clause: PinnedWhereExpr, params: CassandraParams): Array<unknown> | null {
|
||||
if (clause.kind === 'eq') return [getParam(params, clause.param)];
|
||||
const values = getParam(params, clause.param);
|
||||
if (values === null || values === undefined) return [];
|
||||
if (values instanceof Set) return [...values];
|
||||
if (Array.isArray(values)) return [...values];
|
||||
return null;
|
||||
}
|
||||
|
||||
function isKeyComparableValue(value: unknown): boolean {
|
||||
if (value === null || value === undefined) return true;
|
||||
if (typeof value === 'string' || typeof value === 'boolean' || typeof value === 'bigint') return true;
|
||||
if (typeof value === 'number') return Number.isFinite(value);
|
||||
if (Buffer.isBuffer(value)) return true;
|
||||
if (value instanceof Date) return !Number.isNaN(value.getTime());
|
||||
return false;
|
||||
}
|
||||
|
||||
function keySegments(values: ReadonlyArray<unknown>): Array<string> | null {
|
||||
const segments: Array<string> = [];
|
||||
for (const value of values) {
|
||||
let segment: string;
|
||||
try {
|
||||
segment = valueKey(value);
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
if (typeof segment !== 'string') return null;
|
||||
segments.push(segment);
|
||||
}
|
||||
return segments;
|
||||
}
|
||||
|
||||
function combinationCount(segmentLists: ReadonlyArray<ReadonlyArray<string>>): number {
|
||||
let total = 1;
|
||||
for (const segments of segmentLists) total *= segments.length;
|
||||
return total;
|
||||
}
|
||||
|
||||
function keyPrefixes(segmentLists: ReadonlyArray<ReadonlyArray<string>>): Array<string> {
|
||||
let prefixes: Array<string> | null = null;
|
||||
for (const segments of segmentLists) {
|
||||
const next: Array<string> = [];
|
||||
for (const segment of segments) {
|
||||
if (prefixes === null) {
|
||||
next.push(segment);
|
||||
continue;
|
||||
}
|
||||
for (const prefix of prefixes) next.push(`${prefix}${VALUE_SEPARATOR}${segment}`);
|
||||
}
|
||||
prefixes = next;
|
||||
}
|
||||
return [...new Set(prefixes ?? [])];
|
||||
}
|
||||
|
||||
function keyRangeLowerBound(prefix: string): string {
|
||||
return `${prefix}${VALUE_SEPARATOR}`;
|
||||
}
|
||||
|
||||
function keyRangeUpperBound(prefix: string): string {
|
||||
return `${prefix}${KEY_RANGE_UPPER}`;
|
||||
}
|
||||
|
||||
interface PinnedColumn {
|
||||
values: Array<unknown>;
|
||||
segments: Array<string>;
|
||||
}
|
||||
|
||||
function pinnedColumns(clauses: ReadonlyArray<PinnedWhereExpr>, params: CassandraParams): Array<PinnedColumn> {
|
||||
const pinned: Array<PinnedColumn> = [];
|
||||
for (const clause of clauses) {
|
||||
const values = clauseValues(clause, params);
|
||||
if (values === null || values.length === 0) break;
|
||||
const segments = keySegments(values);
|
||||
if (segments === null) break;
|
||||
pinned.push({values, segments});
|
||||
}
|
||||
return pinned;
|
||||
}
|
||||
|
||||
function planIsExact(meta: KvQueryMeta, shape: QueryShape, pinned: ReadonlyArray<PinnedColumn>): boolean {
|
||||
if (pinned.length !== shape.clauseCount) return false;
|
||||
if (meta.limit !== undefined || meta.orderBy !== undefined) return false;
|
||||
for (const column of pinned) {
|
||||
for (const value of column.values) {
|
||||
if (!isKeyComparableValue(value)) return false;
|
||||
}
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
export function buildCandidatePlan(meta: KvQueryMeta, params: CassandraParams): QueryPlan {
|
||||
const shape = queryShape(meta);
|
||||
for (const clause of shape.inClauses) {
|
||||
const values = clauseValues(clause, params);
|
||||
if (values !== null && values.length === 0) {
|
||||
return {candidates: {kind: 'none'}, exact: true};
|
||||
}
|
||||
}
|
||||
const primaryKey = meta.table.primaryKey as ReadonlyArray<string>;
|
||||
const pinned = pinnedColumns(shape.leadingClauses, params);
|
||||
let leading = pinned.length;
|
||||
while (leading > 0) {
|
||||
const segmentLists = pinned.slice(0, leading).map((column) => column.segments);
|
||||
const cap = leading === primaryKey.length ? MAX_ROW_KEY_COMBINATIONS : MAX_PREFIX_RANGES;
|
||||
if (combinationCount(segmentLists) <= cap) break;
|
||||
leading -= 1;
|
||||
}
|
||||
if (leading > 0) {
|
||||
const columns = pinned.slice(0, leading);
|
||||
const prefixes = keyPrefixes(columns.map((column) => column.segments));
|
||||
const exact = planIsExact(meta, shape, columns);
|
||||
if (leading === primaryKey.length) {
|
||||
return {candidates: {kind: 'rowKeys', rowKeys: prefixes}, exact};
|
||||
}
|
||||
if (prefixes.length === 1) {
|
||||
const prefix = prefixes[0]!;
|
||||
return {
|
||||
candidates: {kind: 'range', lowerBound: keyRangeLowerBound(prefix), upperBound: keyRangeUpperBound(prefix)},
|
||||
exact,
|
||||
};
|
||||
}
|
||||
return {
|
||||
candidates: {
|
||||
kind: 'ranges',
|
||||
lowerBounds: prefixes.map(keyRangeLowerBound),
|
||||
upperBounds: prefixes.map(keyRangeUpperBound),
|
||||
},
|
||||
exact,
|
||||
};
|
||||
}
|
||||
if (shape.partitionClauses) {
|
||||
const partition = pinnedColumns(shape.partitionClauses, params);
|
||||
if (
|
||||
partition.length === shape.partitionClauses.length &&
|
||||
combinationCount(partition.map((column) => column.segments)) <= MAX_ROW_KEY_COMBINATIONS
|
||||
) {
|
||||
return {
|
||||
candidates: {kind: 'partitionKeys', partitionKeys: keyPrefixes(partition.map((column) => column.segments))},
|
||||
exact: planIsExact(meta, shape, partition),
|
||||
};
|
||||
}
|
||||
}
|
||||
return {candidates: {kind: 'scan'}, exact: planIsExact(meta, shape, [])};
|
||||
}
|
||||
|
||||
export function planFragments(plan: CandidatePlan): PlanFragments {
|
||||
switch (plan.kind) {
|
||||
case 'rowKeys':
|
||||
return {predicate: ' AND kv.row_key = ANY($2::text[])', params: [plan.rowKeys]};
|
||||
case 'range':
|
||||
return {
|
||||
predicate: ' AND kv.row_key COLLATE "C" >= $2 AND kv.row_key COLLATE "C" < $3',
|
||||
params: [plan.lowerBound, plan.upperBound],
|
||||
};
|
||||
case 'ranges': {
|
||||
const params: Array<unknown> = [];
|
||||
const arms = plan.lowerBounds.map((lowerBound, index) => {
|
||||
params.push(lowerBound, plan.upperBounds[index]);
|
||||
return `(kv.row_key COLLATE "C" >= $${params.length} AND kv.row_key COLLATE "C" < $${params.length + 1})`;
|
||||
});
|
||||
return {predicate: ` AND (${arms.join(' OR ')})`, params};
|
||||
}
|
||||
case 'partitionKeys':
|
||||
return plan.partitionKeys.length === 1
|
||||
? {predicate: ' AND kv.partition_key = $2', params: [plan.partitionKeys[0]]}
|
||||
: {predicate: ' AND kv.partition_key = ANY($2::text[])', params: [plan.partitionKeys]};
|
||||
default:
|
||||
return {predicate: '', params: []};
|
||||
}
|
||||
}
|
||||
|
||||
function logWarn(details: Record<string, unknown>, message: string): void {
|
||||
try {
|
||||
Logger.warn(details, message);
|
||||
} catch {}
|
||||
}
|
||||
|
||||
function logError(details: Record<string, unknown>, message: string): void {
|
||||
try {
|
||||
Logger.error(details, message);
|
||||
} catch {}
|
||||
}
|
||||
|
||||
const fullScanLoggedAt = new Map<string, number>();
|
||||
|
||||
function logFullScan(meta: KvQueryMeta): void {
|
||||
const shape = queryShape(meta);
|
||||
const key = `${meta.table.name}|${meta.action}|${shape.summary}`;
|
||||
const now = Date.now();
|
||||
const last = fullScanLoggedAt.get(key);
|
||||
if (last !== undefined && now - last < FULL_SCAN_LOG_INTERVAL_MS) return;
|
||||
if (fullScanLoggedAt.size >= FULL_SCAN_LOG_KEY_LIMIT) fullScanLoggedAt.clear();
|
||||
fullScanLoggedAt.set(key, now);
|
||||
logWarn({table: meta.table.name, action: meta.action, where: shape.summary || 'none'}, 'Postgres KV full table scan');
|
||||
}
|
||||
|
||||
function ttlExpiresAt(meta: KvQueryMeta, params: CassandraParams): Date | null | undefined {
|
||||
@@ -355,6 +634,24 @@ function parseEqWhere(whereSql: string, cql: string): ReadonlyArray<EqWhereExpr>
|
||||
});
|
||||
}
|
||||
|
||||
async function repairInvalidPostgresKvIndexes(db: PostgresQueryable, kvTable: string): Promise<void> {
|
||||
const invalid = await db.query<{index_name: string; index_def: string}>(
|
||||
`SELECT cls.relname AS index_name, pg_get_indexdef(idx.indexrelid) AS index_def
|
||||
FROM pg_index idx
|
||||
JOIN pg_class cls ON cls.oid = idx.indexrelid
|
||||
WHERE idx.indrelid = to_regclass($1)
|
||||
AND NOT idx.indisvalid
|
||||
AND NOT idx.indisprimary
|
||||
AND NOT idx.indisunique`,
|
||||
[kvTable],
|
||||
);
|
||||
for (const row of invalid.rows) {
|
||||
logError({table: kvTable, index: row.index_name}, 'Postgres KV index is invalid, rebuilding it');
|
||||
await db.query(`DROP INDEX IF EXISTS ${quoteIdentifier(row.index_name)}`);
|
||||
await db.query(row.index_def);
|
||||
}
|
||||
}
|
||||
|
||||
export async function ensurePostgresKvSchema(client: IPostgresClient): Promise<void> {
|
||||
const kvTable = client.kvTable();
|
||||
const table = quoteIdentifier(kvTable);
|
||||
@@ -387,12 +684,22 @@ CREATE TABLE IF NOT EXISTS ${table} (
|
||||
await db.query(
|
||||
`CREATE INDEX IF NOT EXISTS ${quoteIdentifier(`${kvTable}_message_reactions_message_idx`)} ON ${table} (partition_key, ((CASE WHEN row_data -> 'message_id' ->> 'value' ~ '^-?[0-9]+$' THEN (row_data -> 'message_id' ->> 'value')::bigint END))) WHERE table_name = 'message_reactions'`,
|
||||
);
|
||||
await db.query(`
|
||||
await repairInvalidPostgresKvIndexes(db, kvTable);
|
||||
const pending = await db.query(`
|
||||
SELECT 1
|
||||
FROM ${table}
|
||||
WHERE table_name = 'messages'
|
||||
AND partition_key = row_key
|
||||
AND split_part(row_key, chr(31), 3) <> ''
|
||||
LIMIT 1`);
|
||||
if (pending.rows.length > 0) {
|
||||
await db.query(`
|
||||
UPDATE ${table}
|
||||
SET partition_key = split_part(row_key, chr(31), 1) || chr(31) || split_part(row_key, chr(31), 2)
|
||||
WHERE table_name = 'messages'
|
||||
AND partition_key = row_key
|
||||
AND split_part(row_key, chr(31), 3) <> ''`);
|
||||
}
|
||||
await db.query(`DROP INDEX IF EXISTS ${quoteIdentifier(`${kvTable}_partition_idx`)}`);
|
||||
});
|
||||
}
|
||||
@@ -434,9 +741,9 @@ export class PostgresKvQueryExecutor {
|
||||
const meta = this.meta(query);
|
||||
switch (meta.action) {
|
||||
case 'select':
|
||||
return (await this.select(meta, query.params, db)) as Array<T>;
|
||||
return (await this.select(meta, query.params, buildCandidatePlan(meta, query.params), db)) as Array<T>;
|
||||
case 'count':
|
||||
return [{count: (await this.select(meta, query.params, db)).length}] as Array<T>;
|
||||
return (await this.count(meta, query.params, db)) as Array<T>;
|
||||
case 'upsert':
|
||||
return (await this.upsert(meta, query.params, db)) as Array<T>;
|
||||
case 'insert':
|
||||
@@ -495,42 +802,51 @@ export class PostgresKvQueryExecutor {
|
||||
return meta;
|
||||
}
|
||||
|
||||
private async candidates(
|
||||
meta: KvQueryMeta,
|
||||
params: CassandraParams,
|
||||
db: PostgresQueryable,
|
||||
): Promise<Array<StoredRow>> {
|
||||
const rowKeys = fullRowKeysFromWhere(meta, params);
|
||||
if (rowKeys) {
|
||||
const result = await db.query<StoredRow>(
|
||||
`SELECT row_key, row_data FROM ${this.table} WHERE table_name = $1 AND row_key = ANY($2::text[]) AND (expires_at IS NULL OR expires_at > now())`,
|
||||
[meta.table.name, rowKeys],
|
||||
);
|
||||
return result.rows;
|
||||
}
|
||||
if (hasFullPartition(meta)) {
|
||||
const result = await db.query<StoredRow>(
|
||||
`SELECT row_key, row_data FROM ${this.table} WHERE table_name = $1 AND partition_key = $2 AND (expires_at IS NULL OR expires_at > now())`,
|
||||
[meta.table.name, partitionKeyFromParams(meta, params)],
|
||||
);
|
||||
return result.rows;
|
||||
}
|
||||
private async candidates(meta: KvQueryMeta, plan: QueryPlan, db: PostgresQueryable): Promise<Array<StoredRow>> {
|
||||
if (plan.candidates.kind === 'none') return [];
|
||||
if (plan.candidates.kind === 'scan') logFullScan(meta);
|
||||
const fragments = planFragments(plan.candidates);
|
||||
const result = await db.query<StoredRow>(
|
||||
`SELECT row_key, row_data FROM ${this.table} WHERE table_name = $1 AND (expires_at IS NULL OR expires_at > now())`,
|
||||
[meta.table.name],
|
||||
`SELECT kv.row_key, kv.row_data FROM ${this.table} kv WHERE kv.table_name = $1${fragments.predicate} AND (kv.expires_at IS NULL OR kv.expires_at > now())`,
|
||||
[meta.table.name, ...fragments.params],
|
||||
);
|
||||
return result.rows;
|
||||
}
|
||||
|
||||
private async select(meta: KvQueryMeta, params: CassandraParams, db: PostgresQueryable): Promise<Array<Row>> {
|
||||
let rows = (await this.candidates(meta, params, db))
|
||||
.map((stored) => decodeRow(stored.row_data))
|
||||
private matchingRows(meta: KvQueryMeta, stored: ReadonlyArray<StoredRow>, params: CassandraParams): Array<Row> {
|
||||
const required = queryShape(meta).requiredColumns;
|
||||
return stored
|
||||
.map((entry) => (required ? decodeRowColumns(entry.row_data, required) : decodeRow(entry.row_data)))
|
||||
.filter((row) => matchesWhere(row, meta.where as ReadonlyArray<WhereExpr<Row>> | undefined, params));
|
||||
}
|
||||
|
||||
private async select(
|
||||
meta: KvQueryMeta,
|
||||
params: CassandraParams,
|
||||
plan: QueryPlan,
|
||||
db: PostgresQueryable,
|
||||
): Promise<Array<Row>> {
|
||||
let rows = this.matchingRows(meta, await this.candidates(meta, plan, db), params);
|
||||
rows = sortRows(meta, rows);
|
||||
if (typeof meta.limit === 'number') rows = rows.slice(0, meta.limit);
|
||||
return rows.map((row) => projectRow(row, meta.columns as ReadonlyArray<string> | undefined));
|
||||
}
|
||||
|
||||
private async count(meta: KvQueryMeta, params: CassandraParams, db: PostgresQueryable): Promise<Array<Row>> {
|
||||
const plan = buildCandidatePlan(meta, params);
|
||||
if (!plan.exact) {
|
||||
return [{count: (await this.select(meta, params, plan, db)).length}];
|
||||
}
|
||||
if (plan.candidates.kind === 'none') return [{count: 0}];
|
||||
if (plan.candidates.kind === 'scan') logFullScan(meta);
|
||||
const fragments = planFragments(plan.candidates);
|
||||
const result = await db.query<{count: string}>(
|
||||
`SELECT count(*) AS count FROM ${this.table} kv WHERE kv.table_name = $1${fragments.predicate} AND (kv.expires_at IS NULL OR kv.expires_at > now())`,
|
||||
[meta.table.name, ...fragments.params],
|
||||
);
|
||||
return [{count: Number(result.rows[0]?.count ?? 0)}];
|
||||
}
|
||||
|
||||
private async upsert(meta: KvQueryMeta, params: CassandraParams, db: PostgresQueryable): Promise<Array<Row>> {
|
||||
const incoming = rowFromParams(meta, params);
|
||||
const key = rowKey(meta, incoming);
|
||||
@@ -587,10 +903,26 @@ DO UPDATE SET partition_key = EXCLUDED.partition_key, row_data = EXCLUDED.row_da
|
||||
}
|
||||
|
||||
private async delete(meta: KvQueryMeta, params: CassandraParams, db: PostgresQueryable): Promise<void> {
|
||||
const rows = await this.candidates(meta, params, db);
|
||||
const plan = buildCandidatePlan(meta, params);
|
||||
if (plan.candidates.kind === 'none') return;
|
||||
if (plan.exact) {
|
||||
if (plan.candidates.kind === 'scan') logFullScan(meta);
|
||||
const fragments = planFragments(plan.candidates);
|
||||
await db.query(
|
||||
`DELETE FROM ${this.table} kv WHERE kv.table_name = $1${fragments.predicate} AND (kv.expires_at IS NULL OR kv.expires_at > now())`,
|
||||
[meta.table.name, ...fragments.params],
|
||||
);
|
||||
return;
|
||||
}
|
||||
const whereColumns = queryShape(meta).whereColumns;
|
||||
const rows = await this.candidates(meta, plan, db);
|
||||
const matchingKeys = rows
|
||||
.filter((stored) =>
|
||||
matchesWhere(decodeRow(stored.row_data), meta.where as ReadonlyArray<WhereExpr<Row>> | undefined, params),
|
||||
matchesWhere(
|
||||
decodeRowColumns(stored.row_data, whereColumns),
|
||||
meta.where as ReadonlyArray<WhereExpr<Row>> | undefined,
|
||||
params,
|
||||
),
|
||||
)
|
||||
.map((stored) => stored.row_key);
|
||||
if (matchingKeys.length === 0) return;
|
||||
|
||||
Reference in new issue
Block a user