mirror of
https://github.com/fluxerapp/fluxer
synced 2026-10-07 19:22:14 +09:00
perf(kv): collate key columns in C and drop the duplicate index (#2162)
This commit is contained in:
@@ -693,6 +693,19 @@ WHERE idx.indrelid = to_regclass($1)
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async function postgresKvRowKeyIsCCollated(db: PostgresQueryable, kvTable: string): Promise<boolean> {
|
||||||
|
const result = await db.query<{c_collated: boolean}>(
|
||||||
|
`SELECT col.collname = 'C' AND col.collnamespace = 'pg_catalog'::regnamespace AS c_collated
|
||||||
|
FROM pg_attribute att
|
||||||
|
JOIN pg_collation col ON col.oid = att.attcollation
|
||||||
|
WHERE att.attrelid = to_regclass($1)
|
||||||
|
AND att.attname = 'row_key'
|
||||||
|
AND NOT att.attisdropped`,
|
||||||
|
[kvTable],
|
||||||
|
);
|
||||||
|
return result.rows[0]?.c_collated === true;
|
||||||
|
}
|
||||||
|
|
||||||
export async function ensurePostgresKvSchema(client: IPostgresClient): Promise<void> {
|
export async function ensurePostgresKvSchema(client: IPostgresClient): Promise<void> {
|
||||||
const kvTable = client.kvTable();
|
const kvTable = client.kvTable();
|
||||||
const table = quoteIdentifier(kvTable);
|
const table = quoteIdentifier(kvTable);
|
||||||
@@ -703,8 +716,8 @@ export async function ensurePostgresKvSchema(client: IPostgresClient): Promise<v
|
|||||||
await db.query(`
|
await db.query(`
|
||||||
CREATE TABLE IF NOT EXISTS ${table} (
|
CREATE TABLE IF NOT EXISTS ${table} (
|
||||||
table_name text NOT NULL,
|
table_name text NOT NULL,
|
||||||
partition_key text NOT NULL,
|
partition_key text COLLATE "C" NOT NULL,
|
||||||
row_key text NOT NULL,
|
row_key text COLLATE "C" NOT NULL,
|
||||||
row_data jsonb NOT NULL,
|
row_data jsonb NOT NULL,
|
||||||
expires_at timestamptz,
|
expires_at timestamptz,
|
||||||
updated_at timestamptz NOT NULL DEFAULT now(),
|
updated_at timestamptz NOT NULL DEFAULT now(),
|
||||||
@@ -713,9 +726,11 @@ CREATE TABLE IF NOT EXISTS ${table} (
|
|||||||
await db.query(
|
await db.query(
|
||||||
`CREATE INDEX IF NOT EXISTS ${quoteIdentifier(`${kvTable}_partition_row_idx`)} ON ${table} (table_name, partition_key, row_key)`,
|
`CREATE INDEX IF NOT EXISTS ${quoteIdentifier(`${kvTable}_partition_row_idx`)} ON ${table} (table_name, partition_key, row_key)`,
|
||||||
);
|
);
|
||||||
await db.query(
|
if (!(await postgresKvRowKeyIsCCollated(db, kvTable))) {
|
||||||
`CREATE INDEX IF NOT EXISTS ${quoteIdentifier(`${kvTable}_row_key_c_idx`)} ON ${table} (table_name, row_key COLLATE "C")`,
|
await db.query(
|
||||||
);
|
`CREATE INDEX IF NOT EXISTS ${quoteIdentifier(`${kvTable}_row_key_c_idx`)} ON ${table} (table_name, row_key COLLATE "C")`,
|
||||||
|
);
|
||||||
|
}
|
||||||
await db.query(
|
await db.query(
|
||||||
`CREATE INDEX IF NOT EXISTS ${quoteIdentifier(`${kvTable}_expires_idx`)} ON ${table} (expires_at) WHERE expires_at IS NOT NULL`,
|
`CREATE INDEX IF NOT EXISTS ${quoteIdentifier(`${kvTable}_expires_idx`)} ON ${table} (expires_at) WHERE expires_at IS NOT NULL`,
|
||||||
);
|
);
|
||||||
|
|||||||
@@ -315,15 +315,52 @@ suite('postgres kv upgrade safety', () => {
|
|||||||
) d`);
|
) d`);
|
||||||
const schemaDiff = await raw.query<{n: string}>(`
|
const schemaDiff = await raw.query<{n: string}>(`
|
||||||
SELECT count(*) AS n FROM (
|
SELECT count(*) AS n FROM (
|
||||||
(SELECT replace(indexdef, '${OLD}', 'KV') FROM pg_indexes WHERE tablename = '${OLD}'
|
(SELECT replace(indexdef, '${OLD}', 'KV') FROM pg_indexes
|
||||||
|
WHERE tablename = '${OLD}' AND indexname <> '${OLD}_row_key_c_idx'
|
||||||
EXCEPT SELECT replace(indexdef, '${NEW}', 'KV') FROM pg_indexes WHERE tablename = '${NEW}')
|
EXCEPT SELECT replace(indexdef, '${NEW}', 'KV') FROM pg_indexes WHERE tablename = '${NEW}')
|
||||||
UNION ALL
|
UNION ALL
|
||||||
(SELECT replace(indexdef, '${NEW}', 'KV') FROM pg_indexes WHERE tablename = '${NEW}'
|
(SELECT replace(indexdef, '${NEW}', 'KV') FROM pg_indexes WHERE tablename = '${NEW}'
|
||||||
EXCEPT SELECT replace(indexdef, '${OLD}', 'KV') FROM pg_indexes WHERE tablename = '${OLD}')
|
EXCEPT SELECT replace(indexdef, '${OLD}', 'KV') FROM pg_indexes WHERE tablename = '${OLD}')
|
||||||
) d`);
|
) d`);
|
||||||
|
const collations = await raw.query<{tablename: string; attname: string; collname: string}>(`
|
||||||
|
SELECT cls.relname AS tablename, att.attname, col.collname
|
||||||
|
FROM pg_attribute att
|
||||||
|
JOIN pg_class cls ON cls.oid = att.attrelid
|
||||||
|
JOIN pg_collation col ON col.oid = att.attcollation
|
||||||
|
WHERE att.attrelid IN ('${OLD}'::regclass, '${NEW}'::regclass)
|
||||||
|
AND att.attname IN ('partition_key', 'row_key')
|
||||||
|
ORDER BY cls.relname, att.attname`);
|
||||||
|
const cIndexes = await raw.query<{tablename: string}>(
|
||||||
|
`SELECT tablename FROM pg_indexes WHERE indexname IN ('${OLD}_row_key_c_idx', '${NEW}_row_key_c_idx')`,
|
||||||
|
);
|
||||||
console.log(`key diffs=${diff.rows[0]!.n} index-definition diffs=${schemaDiff.rows[0]!.n}`);
|
console.log(`key diffs=${diff.rows[0]!.n} index-definition diffs=${schemaDiff.rows[0]!.n}`);
|
||||||
|
console.log(collations.rows.map((r) => `${r.tablename}.${r.attname}=${r.collname}`).join(' '));
|
||||||
expect(Number(diff.rows[0]!.n)).toBe(0);
|
expect(Number(diff.rows[0]!.n)).toBe(0);
|
||||||
expect(Number(schemaDiff.rows[0]!.n)).toBe(0);
|
expect(Number(schemaDiff.rows[0]!.n)).toBe(0);
|
||||||
|
expect(collations.rows.filter((r) => r.tablename === OLD).map((r) => r.collname)).toEqual(['default', 'default']);
|
||||||
|
expect(collations.rows.filter((r) => r.tablename === NEW).map((r) => r.collname)).toEqual(['C', 'C']);
|
||||||
|
expect(cIndexes.rows.map((r) => r.tablename)).toEqual([OLD]);
|
||||||
|
});
|
||||||
|
|
||||||
|
it('never recollates an existing table and keeps its C index', async () => {
|
||||||
|
const OLD = `${KV}_keep_old`;
|
||||||
|
await raw.query(`DROP TABLE IF EXISTS ${OLD}`);
|
||||||
|
const oldClient = new TableClient(raw, OLD);
|
||||||
|
await legacyEnsurePostgresKvSchema(oldClient);
|
||||||
|
await ensurePostgresKvSchema(oldClient);
|
||||||
|
await ensurePostgresKvSchema(oldClient);
|
||||||
|
const collations = await raw.query<{attname: string; collname: string}>(`
|
||||||
|
SELECT att.attname, col.collname
|
||||||
|
FROM pg_attribute att
|
||||||
|
JOIN pg_collation col ON col.oid = att.attcollation
|
||||||
|
WHERE att.attrelid = '${OLD}'::regclass AND att.attname IN ('partition_key', 'row_key')
|
||||||
|
ORDER BY att.attname`);
|
||||||
|
const indexes = await raw.query<{indexname: string}>(
|
||||||
|
`SELECT indexname FROM pg_indexes WHERE tablename = '${OLD}' ORDER BY indexname`,
|
||||||
|
);
|
||||||
|
console.log(`after boot on a legacy table: ${indexes.rows.map((r) => r.indexname).join(', ')}`);
|
||||||
|
expect(collations.rows.map((r) => r.collname)).toEqual(['default', 'default']);
|
||||||
|
expect(indexes.rows.map((r) => r.indexname)).toContain(`${OLD}_row_key_c_idx`);
|
||||||
});
|
});
|
||||||
|
|
||||||
it('survives three simultaneous boots (two api replicas and a worker)', async () => {
|
it('survives three simultaneous boots (two api replicas and a worker)', async () => {
|
||||||
|
|||||||
@@ -3,7 +3,7 @@
|
|||||||
use crate::config::ServiceConfig;
|
use crate::config::ServiceConfig;
|
||||||
use anyhow::Context;
|
use anyhow::Context;
|
||||||
use chrono::{DateTime, Utc};
|
use chrono::{DateTime, Utc};
|
||||||
use deadpool_postgres::{Client, Manager, Pool, Runtime};
|
use deadpool_postgres::{Client, Manager, Pool, Runtime, Transaction};
|
||||||
use rustls::RootCertStore;
|
use rustls::RootCertStore;
|
||||||
use serde_json::{Map, Number, Value};
|
use serde_json::{Map, Number, Value};
|
||||||
use std::io::Cursor;
|
use std::io::Cursor;
|
||||||
@@ -129,6 +129,25 @@ fn build_disabled_tls_connector() -> MakeRustlsConnect {
|
|||||||
MakeRustlsConnect::new(tls_config)
|
MakeRustlsConnect::new(tls_config)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn row_key_is_c_collated(
|
||||||
|
transaction: &Transaction<'_>,
|
||||||
|
kv_table: &str,
|
||||||
|
) -> anyhow::Result<bool> {
|
||||||
|
let row = transaction
|
||||||
|
.query_opt(
|
||||||
|
r#"SELECT col.collname = 'C' AND col.collnamespace = 'pg_catalog'::regnamespace AS c_collated
|
||||||
|
FROM pg_attribute att
|
||||||
|
JOIN pg_collation col ON col.oid = att.attcollation
|
||||||
|
WHERE att.attrelid = to_regclass($1)
|
||||||
|
AND att.attname = 'row_key'
|
||||||
|
AND NOT att.attisdropped"#,
|
||||||
|
&[&kv_table],
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.context("failed to inspect Postgres KV row_key collation")?;
|
||||||
|
Ok(row.and_then(|row| row.get::<_, Option<bool>>("c_collated")) == Some(true))
|
||||||
|
}
|
||||||
|
|
||||||
pub async fn ensure_kv_schema(pool: &Pool, kv_table: &str) -> anyhow::Result<()> {
|
pub async fn ensure_kv_schema(pool: &Pool, kv_table: &str) -> anyhow::Result<()> {
|
||||||
let table = quote_identifier(kv_table)?;
|
let table = quote_identifier(kv_table)?;
|
||||||
let old_partition_index = quote_identifier(&format!("{kv_table}_partition_idx"))?;
|
let old_partition_index = quote_identifier(&format!("{kv_table}_partition_idx"))?;
|
||||||
@@ -166,15 +185,29 @@ pub async fn ensure_kv_schema(pool: &Pool, kv_table: &str) -> anyhow::Result<()>
|
|||||||
r#"
|
r#"
|
||||||
CREATE TABLE IF NOT EXISTS {table} (
|
CREATE TABLE IF NOT EXISTS {table} (
|
||||||
table_name text NOT NULL,
|
table_name text NOT NULL,
|
||||||
partition_key text NOT NULL,
|
partition_key text COLLATE "C" NOT NULL,
|
||||||
row_key text NOT NULL,
|
row_key text COLLATE "C" NOT NULL,
|
||||||
row_data jsonb NOT NULL,
|
row_data jsonb NOT NULL,
|
||||||
expires_at timestamptz,
|
expires_at timestamptz,
|
||||||
updated_at timestamptz NOT NULL DEFAULT now(),
|
updated_at timestamptz NOT NULL DEFAULT now(),
|
||||||
PRIMARY KEY (table_name, row_key)
|
PRIMARY KEY (table_name, row_key)
|
||||||
);
|
);
|
||||||
CREATE INDEX IF NOT EXISTS {partition_row_index} ON {table} (table_name, partition_key, row_key);
|
CREATE INDEX IF NOT EXISTS {partition_row_index} ON {table} (table_name, partition_key, row_key);
|
||||||
CREATE INDEX IF NOT EXISTS {row_key_c_index} ON {table} (table_name, row_key COLLATE "C");
|
"#
|
||||||
|
))
|
||||||
|
.await
|
||||||
|
.context("failed to ensure Postgres KV schema")?;
|
||||||
|
if !row_key_is_c_collated(&transaction, kv_table).await? {
|
||||||
|
transaction
|
||||||
|
.batch_execute(&format!(
|
||||||
|
r#"CREATE INDEX IF NOT EXISTS {row_key_c_index} ON {table} (table_name, row_key COLLATE "C");"#
|
||||||
|
))
|
||||||
|
.await
|
||||||
|
.context("failed to ensure Postgres KV schema")?;
|
||||||
|
}
|
||||||
|
transaction
|
||||||
|
.batch_execute(&format!(
|
||||||
|
r#"
|
||||||
CREATE INDEX IF NOT EXISTS {expires_index} ON {table} (expires_at) WHERE expires_at IS NOT NULL;
|
CREATE INDEX IF NOT EXISTS {expires_index} ON {table} (expires_at) WHERE expires_at IS NOT NULL;
|
||||||
CREATE INDEX IF NOT EXISTS {messages_message_index} ON {table} (partition_key, ((CASE WHEN row_data -> 'message_id' ->> 'value' ~ '^-?[0-9]+$' THEN (row_data -> 'message_id' ->> 'value')::bigint END))) WHERE table_name = 'messages';
|
CREATE INDEX IF NOT EXISTS {messages_message_index} ON {table} (partition_key, ((CASE WHEN row_data -> 'message_id' ->> 'value' ~ '^-?[0-9]+$' THEN (row_data -> 'message_id' ->> 'value')::bigint END))) WHERE table_name = 'messages';
|
||||||
CREATE INDEX IF NOT EXISTS {message_reactions_message_index} 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';
|
CREATE INDEX IF NOT EXISTS {message_reactions_message_index} 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';
|
||||||
|
|||||||
Reference in New Issue
Block a user