mirror of
https://github.com/hcengineering/platform.git
synced 2026-09-30 21:45:01 +02:00
UBERF-8595: Fix backup/restore performance (#7188)
Signed-off-by: Andrey Sobolev <haiodo@gmail.com>
This commit is contained in:
+153
-109
@@ -60,12 +60,11 @@ import {
|
||||
type DomainHelperOperations,
|
||||
estimateDocSize,
|
||||
type ServerFindOptions,
|
||||
type TxAdapter,
|
||||
updateHashForDoc
|
||||
toDocInfo,
|
||||
type TxAdapter
|
||||
} from '@hcengineering/server-core'
|
||||
import { createHash } from 'crypto'
|
||||
import type postgres from 'postgres'
|
||||
import { getDocFieldsByDomains, translateDomain } from './schemas'
|
||||
import { getDocFieldsByDomains, getSchema, translateDomain } from './schemas'
|
||||
import { type ValueType } from './types'
|
||||
import {
|
||||
convertDoc,
|
||||
@@ -173,6 +172,7 @@ abstract class PostgresAdapterBase implements DbAdapter {
|
||||
query: DocumentQuery<T>,
|
||||
options?: Pick<FindOptions<T>, 'sort' | 'limit' | 'projection'>
|
||||
): Promise<Iterator<T>> {
|
||||
const schema = getSchema(_domain)
|
||||
const client = await this.client.reserve()
|
||||
let closed = false
|
||||
const cursorName = `cursor_${translateDomain(this.workspaceId.name)}_${translateDomain(_domain)}_${generateId()}`
|
||||
@@ -209,7 +209,7 @@ abstract class PostgresAdapterBase implements DbAdapter {
|
||||
await close(cursorName)
|
||||
return null
|
||||
}
|
||||
return result.map((p) => parseDoc(p as any, _domain))
|
||||
return result.map((p) => parseDoc(p as any, schema))
|
||||
}
|
||||
|
||||
await init()
|
||||
@@ -306,12 +306,14 @@ abstract class PostgresAdapterBase implements DbAdapter {
|
||||
const res = await client.unsafe(
|
||||
`SELECT * FROM ${translateDomain(domain)} WHERE ${translatedQuery} FOR UPDATE`
|
||||
)
|
||||
const docs = res.map((p) => parseDoc(p as any, domain))
|
||||
const schema = getSchema(domain)
|
||||
const docs = res.map((p) => parseDoc(p as any, schema))
|
||||
const domainFields = new Set(getDocFieldsByDomains(domain))
|
||||
for (const doc of docs) {
|
||||
if (doc === undefined) continue
|
||||
const prevAttachedTo = (doc as any).attachedTo
|
||||
TxProcessor.applyUpdate(doc, operations)
|
||||
const converted = convertDoc(domain, doc, this.workspaceId.name)
|
||||
const converted = convertDoc(domain, doc, this.workspaceId.name, domainFields)
|
||||
const params: any[] = [doc._id, this.workspaceId.name]
|
||||
let paramsIndex = params.length + 1
|
||||
const updates: string[] = []
|
||||
@@ -516,6 +518,7 @@ abstract class PostgresAdapterBase implements DbAdapter {
|
||||
let joinIndex: number | undefined
|
||||
let skip = false
|
||||
try {
|
||||
const schema = getSchema(domain)
|
||||
for (const column in row) {
|
||||
if (column.startsWith('reverse_lookup_')) {
|
||||
if (row[column] != null) {
|
||||
@@ -527,7 +530,7 @@ abstract class PostgresAdapterBase implements DbAdapter {
|
||||
if (res === undefined) continue
|
||||
const { obj, key } = res
|
||||
|
||||
const parsed = row[column].map((p: any) => parseDoc(p, domain))
|
||||
const parsed = row[column].map((p: any) => parseDoc(p, schema))
|
||||
obj[key] = parsed
|
||||
}
|
||||
} else if (column.startsWith('lookup_')) {
|
||||
@@ -1112,115 +1115,127 @@ abstract class PostgresAdapterBase implements DbAdapter {
|
||||
find (_ctx: MeasureContext, domain: Domain, recheck?: boolean): StorageIterator {
|
||||
const ctx = _ctx.newChild('find', { domain })
|
||||
|
||||
const getCursorName = (): string => {
|
||||
return `cursor_${translateDomain(this.workspaceId.name)}_${translateDomain(domain)}_${mode}`
|
||||
}
|
||||
|
||||
let initialized: boolean = false
|
||||
let client: postgres.ReservedSql
|
||||
let mode: 'hashed' | 'non_hashed' = 'hashed'
|
||||
let cursorName = getCursorName()
|
||||
const bulkUpdate = new Map<Ref<Doc>, string>()
|
||||
|
||||
const close = async (cursorName: string): Promise<void> => {
|
||||
try {
|
||||
await client.unsafe(`CLOSE ${cursorName}`)
|
||||
await client.unsafe('COMMIT')
|
||||
} catch (err) {
|
||||
ctx.error('Error while closing cursor', { cursorName, err })
|
||||
} finally {
|
||||
client.release()
|
||||
}
|
||||
}
|
||||
|
||||
const init = async (projection: string, query: string): Promise<void> => {
|
||||
cursorName = getCursorName()
|
||||
client = await this.client.reserve()
|
||||
await client.unsafe('BEGIN')
|
||||
await client.unsafe(
|
||||
`DECLARE ${cursorName} CURSOR FOR SELECT ${projection} FROM ${translateDomain(domain)} WHERE "workspaceId" = $1 AND ${query}`,
|
||||
[this.workspaceId.name]
|
||||
)
|
||||
}
|
||||
|
||||
const next = async (limit: number): Promise<Doc[]> => {
|
||||
const result = await client.unsafe(`FETCH ${limit} FROM ${cursorName}`)
|
||||
if (result.length === 0) {
|
||||
return []
|
||||
}
|
||||
return result.filter((it) => it != null).map((it) => parseDoc(it as any, domain))
|
||||
}
|
||||
const tdomain = translateDomain(domain)
|
||||
const schema = getSchema(domain)
|
||||
|
||||
const flush = async (flush = false): Promise<void> => {
|
||||
if (bulkUpdate.size > 1000 || flush) {
|
||||
if (bulkUpdate.size > 0) {
|
||||
await ctx.with('bulk-write-find', {}, () => {
|
||||
const updates = new Map(Array.from(bulkUpdate.entries()).map((it) => [it[0], { '%hash%': it[1] }]))
|
||||
return this.update(ctx, domain, updates)
|
||||
})
|
||||
const entries = Array.from(bulkUpdate.entries())
|
||||
bulkUpdate.clear()
|
||||
const updateClient = await this.client.reserve()
|
||||
try {
|
||||
while (entries.length > 0) {
|
||||
const part = entries.splice(0, 200)
|
||||
const data: string[] = part.flat()
|
||||
const indexes = part.map((val, idx) => `($${2 * idx + 1}::text, $${2 * idx + 2}::text)`).join(', ')
|
||||
await ctx.with('bulk-write-find', {}, () => {
|
||||
return this.retryTxn(updateClient, (client) =>
|
||||
client.unsafe(
|
||||
`
|
||||
UPDATE ${tdomain} SET "%hash%" = update_data.hash
|
||||
FROM (values ${indexes}) AS update_data(_id, hash)
|
||||
WHERE ${tdomain}."workspaceId" = '${this.workspaceId.name}' AND ${tdomain}."_id" = update_data._id
|
||||
`,
|
||||
data
|
||||
)
|
||||
)
|
||||
})
|
||||
}
|
||||
} catch (err: any) {
|
||||
ctx.error('failed to update hash', { err })
|
||||
} finally {
|
||||
updateClient.release()
|
||||
}
|
||||
}
|
||||
bulkUpdate.clear()
|
||||
}
|
||||
}
|
||||
|
||||
const workspaceId = this.workspaceId
|
||||
|
||||
async function * createBulk (projection: string, query: string, limit = 50): AsyncGenerator<Doc[]> {
|
||||
const cursor = client
|
||||
.unsafe(`SELECT ${projection} FROM ${tdomain} WHERE "workspaceId" = '${workspaceId.name}' AND ${query}`)
|
||||
.cursor(limit)
|
||||
try {
|
||||
for await (const part of cursor) {
|
||||
yield part.filter((it) => it != null).map((it) => parseDoc(it as any, schema))
|
||||
}
|
||||
} catch (err: any) {
|
||||
ctx.error('failed to recieve data', { err })
|
||||
}
|
||||
}
|
||||
let bulk: AsyncGenerator<Doc[]>
|
||||
let forcedRecheck = false
|
||||
|
||||
return {
|
||||
next: async () => {
|
||||
if (!initialized) {
|
||||
if (recheck === true) {
|
||||
await this.retryTxn(client, async (client) => {
|
||||
await client`UPDATE ${client(translateDomain(domain))} SET '%hash%' = NULL WHERE "workspaceId" = ${this.workspaceId.name} AND '%hash%' IS NOT NULL`
|
||||
})
|
||||
if (client === undefined) {
|
||||
client = await this.client.reserve()
|
||||
}
|
||||
await init('_id, data', "'%hash%' IS NOT NULL AND '%hash%' <> ''")
|
||||
|
||||
if (recheck === true) {
|
||||
await this.retryTxn(
|
||||
client,
|
||||
(client) =>
|
||||
client`UPDATE ${client(tdomain)} SET "%hash%" = NULL WHERE "workspaceId" = ${this.workspaceId.name} AND "%hash%" IS NOT NULL`
|
||||
)
|
||||
}
|
||||
|
||||
initialized = true
|
||||
await flush(true) // We need to flush, so wrong id documents will be updated.
|
||||
bulk = createBulk('_id, "%hash%"', '"%hash%" IS NOT NULL AND "%hash%" <> \'\'')
|
||||
// bulk = createBulk('_id, "%hash%, data', '"%hash%" IS NOT NULL AND "%hash%" <> \'\'')
|
||||
}
|
||||
let docs = await ctx.with('next', { mode }, () => next(50))
|
||||
if (docs.length === 0 && mode === 'hashed') {
|
||||
await close(cursorName)
|
||||
|
||||
let docs = await ctx.with('next', { mode }, () => bulk.next())
|
||||
|
||||
if (!forcedRecheck && docs.done !== true && docs.value?.length > 0) {
|
||||
// Check if we have wrong hash stored, and update all of them.
|
||||
forcedRecheck = true
|
||||
|
||||
for (const d of docs.value) {
|
||||
const digest: string | null = (d as any)['%hash%']
|
||||
|
||||
const pos = (digest ?? '').indexOf('|')
|
||||
if (pos === -1) {
|
||||
await bulk.return([]) // We need to close generator
|
||||
docs = { done: true, value: undefined }
|
||||
await this.retryTxn(
|
||||
client,
|
||||
(client) =>
|
||||
client`UPDATE ${client(tdomain)} SET "%hash%" = NULL WHERE "workspaceId" = ${this.workspaceId.name} AND "%hash%" IS NOT NULL`
|
||||
)
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if ((docs.done === true || docs.value.length === 0) && mode === 'hashed') {
|
||||
forcedRecheck = true
|
||||
mode = 'non_hashed'
|
||||
await init('*', "'%hash%' IS NULL OR '%hash%' = ''")
|
||||
docs = await ctx.with('next', { mode }, () => next(50))
|
||||
bulk = createBulk('*', '"%hash%" IS NULL OR "%hash%" = \'\'')
|
||||
docs = await ctx.with('next', { mode }, () => bulk.next())
|
||||
}
|
||||
if (docs.length === 0) {
|
||||
if (docs.done === true || docs.value.length === 0) {
|
||||
return []
|
||||
}
|
||||
const result: DocInfo[] = []
|
||||
for (const d of docs) {
|
||||
let digest: string | null = (d as any)['%hash%']
|
||||
if ('%hash%' in d) {
|
||||
delete d['%hash%']
|
||||
}
|
||||
const pos = (digest ?? '').indexOf('|')
|
||||
if (digest == null || digest === '') {
|
||||
const cs = ctx.newChild('calc-size', {})
|
||||
const size = estimateDocSize(d)
|
||||
cs.end()
|
||||
|
||||
const hash = createHash('sha256')
|
||||
updateHashForDoc(hash, d)
|
||||
digest = hash.digest('base64')
|
||||
|
||||
bulkUpdate.set(d._id, `${digest}|${size.toString(16)}`)
|
||||
|
||||
await ctx.with('flush', {}, () => flush())
|
||||
result.push({
|
||||
id: d._id,
|
||||
hash: digest,
|
||||
size
|
||||
})
|
||||
} else {
|
||||
result.push({
|
||||
id: d._id,
|
||||
hash: digest.slice(0, pos),
|
||||
size: parseInt(digest.slice(pos + 1), 16)
|
||||
})
|
||||
}
|
||||
for (const d of docs.value) {
|
||||
result.push(toDocInfo(d, bulkUpdate))
|
||||
}
|
||||
await ctx.with('flush', {}, () => flush())
|
||||
return result
|
||||
},
|
||||
close: async () => {
|
||||
await ctx.with('flush', {}, () => flush(true))
|
||||
await close(cursorName)
|
||||
client?.release()
|
||||
ctx.end()
|
||||
}
|
||||
}
|
||||
@@ -1231,16 +1246,16 @@ abstract class PostgresAdapterBase implements DbAdapter {
|
||||
if (docs.length === 0) {
|
||||
return []
|
||||
}
|
||||
const connection = (await this.getConnection(ctx)) ?? this.client
|
||||
const res =
|
||||
await connection`SELECT * FROM ${connection(translateDomain(domain))} WHERE _id = ANY(${docs}) AND "workspaceId" = ${this.workspaceId.name}`
|
||||
return res.map((p) => parseDocWithProjection(p as any, domain))
|
||||
return await this.withConnection(ctx, async (connection) => {
|
||||
const res =
|
||||
await connection`SELECT * FROM ${connection(translateDomain(domain))} WHERE _id = ANY(${docs}) AND "workspaceId" = ${this.workspaceId.name}`
|
||||
return res.map((p) => parseDocWithProjection(p as any, domain))
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
upload (ctx: MeasureContext, domain: Domain, docs: Doc[]): Promise<void> {
|
||||
return ctx.with('upload', { domain }, async (ctx) => {
|
||||
const arr = docs.concat()
|
||||
const fields = getDocFieldsByDomains(domain)
|
||||
const filedsWithData = [...fields, 'data']
|
||||
const insertFields: string[] = []
|
||||
@@ -1252,15 +1267,26 @@ abstract class PostgresAdapterBase implements DbAdapter {
|
||||
const insertStr = insertFields.join(', ')
|
||||
const onConflictStr = onConflict.join(', ')
|
||||
await this.withConnection(ctx, async (connection) => {
|
||||
while (arr.length > 0) {
|
||||
const part = arr.splice(0, 500)
|
||||
const domainFields = new Set(getDocFieldsByDomains(domain))
|
||||
const toUpload = [...docs]
|
||||
const tdomain = translateDomain(domain)
|
||||
while (toUpload.length > 0) {
|
||||
const part = toUpload.splice(0, 200)
|
||||
const values: any[] = []
|
||||
const vars: string[] = []
|
||||
let index = 1
|
||||
for (let i = 0; i < part.length; i++) {
|
||||
const doc = part[i]
|
||||
const variables: string[] = []
|
||||
const d = convertDoc(domain, doc, this.workspaceId.name)
|
||||
|
||||
const digest: string | null = (doc as any)['%hash%']
|
||||
if ('%hash%' in doc) {
|
||||
delete doc['%hash%']
|
||||
}
|
||||
const size = digest != null ? estimateDocSize(doc) : 0
|
||||
;(doc as any)['%hash%'] = digest == null ? null : `${digest}|${size.toString(16)}`
|
||||
const d = convertDoc(domain, doc, this.workspaceId.name, domainFields)
|
||||
|
||||
values.push(d.workspaceId)
|
||||
variables.push(`$${index++}`)
|
||||
for (const field of fields) {
|
||||
@@ -1273,21 +1299,36 @@ abstract class PostgresAdapterBase implements DbAdapter {
|
||||
}
|
||||
|
||||
const vals = vars.join(',')
|
||||
await this.retryTxn(connection, async (client) => {
|
||||
await client.unsafe(
|
||||
`INSERT INTO ${translateDomain(domain)} ("workspaceId", ${insertStr}) VALUES ${vals}
|
||||
await this.retryTxn(connection, (client) =>
|
||||
client.unsafe(
|
||||
`INSERT INTO ${tdomain} ("workspaceId", ${insertStr}) VALUES ${vals}
|
||||
ON CONFLICT ("workspaceId", _id) DO UPDATE SET ${onConflictStr};`,
|
||||
values
|
||||
)
|
||||
})
|
||||
)
|
||||
}
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
async clean (ctx: MeasureContext, domain: Domain, docs: Ref<Doc>[]): Promise<void> {
|
||||
const connection = (await this.getConnection(ctx)) ?? this.client
|
||||
await connection`DELETE FROM ${connection(translateDomain(domain))} WHERE _id = ANY(${docs}) AND "workspaceId" = ${this.workspaceId.name}`
|
||||
const updateClient = await this.client.reserve()
|
||||
try {
|
||||
const tdomain = translateDomain(domain)
|
||||
const toClean = [...docs]
|
||||
while (toClean.length > 0) {
|
||||
const part = toClean.splice(0, 200)
|
||||
await ctx.with('clean', {}, () => {
|
||||
return this.retryTxn(
|
||||
updateClient,
|
||||
(client) =>
|
||||
client`DELETE FROM ${client(tdomain)} WHERE _id = ANY(${part}) AND "workspaceId" = ${this.workspaceId.name}`
|
||||
)
|
||||
})
|
||||
}
|
||||
} finally {
|
||||
updateClient.release()
|
||||
}
|
||||
}
|
||||
|
||||
groupBy<T, P extends Doc>(
|
||||
@@ -1318,8 +1359,10 @@ abstract class PostgresAdapterBase implements DbAdapter {
|
||||
try {
|
||||
const res =
|
||||
await client`SELECT * FROM ${client(translateDomain(domain))} WHERE _id = ANY(${ids}) AND "workspaceId" = ${this.workspaceId.name} FOR UPDATE`
|
||||
const docs = res.map((p) => parseDoc(p as any, domain))
|
||||
const schema = getSchema(domain)
|
||||
const docs = res.map((p) => parseDoc(p as any, schema))
|
||||
const map = new Map(docs.map((d) => [d._id, d]))
|
||||
const domainFields = new Set(getDocFieldsByDomains(domain))
|
||||
for (const [_id, ops] of operations) {
|
||||
const doc = map.get(_id)
|
||||
if (doc === undefined) continue
|
||||
@@ -1328,7 +1371,7 @@ abstract class PostgresAdapterBase implements DbAdapter {
|
||||
;(op as any)['%hash%'] = null
|
||||
}
|
||||
TxProcessor.applyUpdate(doc, op)
|
||||
const converted = convertDoc(domain, doc, this.workspaceId.name)
|
||||
const converted = convertDoc(domain, doc, this.workspaceId.name, domainFields)
|
||||
|
||||
const columns: string[] = []
|
||||
const { extractedFields, remainingData } = parseUpdate(domain, op)
|
||||
@@ -1361,12 +1404,13 @@ abstract class PostgresAdapterBase implements DbAdapter {
|
||||
columns.push(field)
|
||||
}
|
||||
await this.withConnection(ctx, async (connection) => {
|
||||
const domainFields = new Set(getDocFieldsByDomains(domain))
|
||||
while (docs.length > 0) {
|
||||
const part = docs.splice(0, 500)
|
||||
const values: DBDoc[] = []
|
||||
for (let i = 0; i < part.length; i++) {
|
||||
const doc = part[i]
|
||||
const d = convertDoc(domain, doc, this.workspaceId.name)
|
||||
const d = convertDoc(domain, doc, this.workspaceId.name, domainFields)
|
||||
values.push(d)
|
||||
}
|
||||
await this.retryTxn(connection, async (client) => {
|
||||
@@ -1559,14 +1603,14 @@ class PostgresAdapter extends PostgresAdapterBase {
|
||||
_id: Ref<Doc>,
|
||||
forUpdate: boolean = false
|
||||
): Promise<Doc | undefined> {
|
||||
const domain = this.hierarchy.getDomain(_class)
|
||||
return ctx.with('find-doc', { _class }, async () => {
|
||||
const res =
|
||||
await client`SELECT * FROM ${this.client(translateDomain(this.hierarchy.getDomain(_class)))} WHERE _id = ${_id} AND "workspaceId" = ${this.workspaceId.name} ${
|
||||
await client`SELECT * FROM ${this.client(translateDomain(domain))} WHERE _id = ${_id} AND "workspaceId" = ${this.workspaceId.name} ${
|
||||
forUpdate ? client` FOR UPDATE` : client``
|
||||
}`
|
||||
const dbDoc = res[0]
|
||||
const domain = this.hierarchy.getDomain(_class)
|
||||
return dbDoc !== undefined ? parseDoc(dbDoc as any, domain) : undefined
|
||||
return dbDoc !== undefined ? parseDoc(dbDoc as any, getSchema(domain)) : undefined
|
||||
})
|
||||
}
|
||||
|
||||
@@ -1622,7 +1666,7 @@ class PostgresTxAdapter extends PostgresAdapterBase implements TxAdapter {
|
||||
const res = await this
|
||||
.client`SELECT * FROM ${this.client(translateDomain(DOMAIN_MODEL_TX))} WHERE "workspaceId" = ${this.workspaceId.name} ORDER BY _id ASC, "modifiedOn" ASC`
|
||||
|
||||
const model = res.map((p) => parseDoc<Tx>(p as any, DOMAIN_MODEL_TX))
|
||||
const model = res.map((p) => parseDoc<Tx>(p as any, getSchema(DOMAIN_MODEL_TX)))
|
||||
// We need to put all core.account.System transactions first
|
||||
const systemTx: Tx[] = []
|
||||
const userTx: Tx[] = []
|
||||
|
||||
@@ -195,7 +195,6 @@ class PostgresClientReferenceImpl {
|
||||
this.onclose()
|
||||
const cl = await this.client
|
||||
await cl.end()
|
||||
console.log('Closed postgres connection')
|
||||
})()
|
||||
}
|
||||
}
|
||||
@@ -261,7 +260,12 @@ export function getDBClient (connectionString: string, database?: string): Postg
|
||||
return new ClientRef(existing)
|
||||
}
|
||||
|
||||
export function convertDoc<T extends Doc> (domain: string, doc: T, workspaceId: string): DBDoc {
|
||||
export function convertDoc<T extends Doc> (
|
||||
domain: string,
|
||||
doc: T,
|
||||
workspaceId: string,
|
||||
domainFields?: Set<string>
|
||||
): DBDoc {
|
||||
const extractedFields: Doc & Record<string, any> = {
|
||||
_id: doc._id,
|
||||
space: doc.space,
|
||||
@@ -273,9 +277,15 @@ export function convertDoc<T extends Doc> (domain: string, doc: T, workspaceId:
|
||||
}
|
||||
const remainingData: Partial<T> = {}
|
||||
|
||||
const extractedFieldsKeys = new Set(Object.keys(extractedFields))
|
||||
|
||||
domainFields = domainFields ?? new Set(getDocFieldsByDomains(domain))
|
||||
|
||||
for (const key in doc) {
|
||||
if (Object.keys(extractedFields).includes(key)) continue
|
||||
if (getDocFieldsByDomains(domain).includes(key)) {
|
||||
if (extractedFieldsKeys.has(key)) {
|
||||
continue
|
||||
}
|
||||
if (domainFields.has(key)) {
|
||||
extractedFields[key] = doc[key]
|
||||
} else {
|
||||
remainingData[key] = doc[key]
|
||||
@@ -432,8 +442,7 @@ export function parseDocWithProjection<T extends Doc> (
|
||||
return res
|
||||
}
|
||||
|
||||
export function parseDoc<T extends Doc> (doc: DBDoc, domain: string): T {
|
||||
const schema = getSchema(domain)
|
||||
export function parseDoc<T extends Doc> (doc: DBDoc, schema: Schema): T {
|
||||
const { workspaceId, data, ...rest } = doc
|
||||
for (const key in rest) {
|
||||
if ((rest as any)[key] === 'NULL' || (rest as any)[key] === null) {
|
||||
|
||||
Reference in New Issue
Block a user