mirror of
https://github.com/fluxerapp/fluxer
synced 2026-10-07 19:22:14 +09:00
fix(cassandra): bound driver in-flight requests per connection (#2122)
This commit is contained in:
@@ -6,6 +6,9 @@ import cassandra from 'cassandra-driver';
|
||||
|
||||
const distance = cassandra.types.distance;
|
||||
|
||||
const MAX_REQUESTS_PER_CONNECTION = 2048;
|
||||
const READ_TIMEOUT_MS = 12000;
|
||||
|
||||
interface CassandraConfig {
|
||||
hosts: Array<string>;
|
||||
port?: number | undefined;
|
||||
@@ -83,12 +86,15 @@ class CassandraClient implements ICassandraClient {
|
||||
port: this.config.port ?? 9042,
|
||||
},
|
||||
pooling: {
|
||||
maxRequestsPerConnection: 32768,
|
||||
maxRequestsPerConnection: MAX_REQUESTS_PER_CONNECTION,
|
||||
coreConnectionsPerHost: {
|
||||
[distance.local]: 4,
|
||||
[distance.remote]: 2,
|
||||
},
|
||||
},
|
||||
socketOptions: {
|
||||
readTimeout: READ_TIMEOUT_MS,
|
||||
},
|
||||
encoding: {
|
||||
map: Map,
|
||||
set: Set,
|
||||
|
||||
@@ -0,0 +1,36 @@
|
||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
|
||||
import {ServiceUnavailableError} from '@fluxer/errors/src/domains/core/ServiceUnavailableError';
|
||||
import cassandra from 'cassandra-driver';
|
||||
import {describe, expect, it} from 'vitest';
|
||||
import {mapCassandraDriverError} from './CassandraQueryExecution';
|
||||
|
||||
function busyConnectionError(): cassandra.errors.BusyConnectionError {
|
||||
return new cassandra.errors.BusyConnectionError('127.0.0.1:9042', 2048, 4);
|
||||
}
|
||||
|
||||
describe('mapCassandraDriverError', () => {
|
||||
it('sheds a busy connection error as a 503', () => {
|
||||
const mapped = mapCassandraDriverError(busyConnectionError());
|
||||
expect(mapped).toBeInstanceOf(ServiceUnavailableError);
|
||||
expect((mapped as ServiceUnavailableError).status).toBe(503);
|
||||
});
|
||||
|
||||
it('sheds a busy connection error nested in a no host available error as a 503', () => {
|
||||
const mapped = mapCassandraDriverError(
|
||||
new cassandra.errors.NoHostAvailableError({'127.0.0.1:9042': busyConnectionError()}),
|
||||
);
|
||||
expect(mapped).toBeInstanceOf(ServiceUnavailableError);
|
||||
expect((mapped as ServiceUnavailableError).status).toBe(503);
|
||||
});
|
||||
|
||||
it('returns other no host available errors unchanged', () => {
|
||||
const err = new cassandra.errors.NoHostAvailableError({'127.0.0.1:9042': new Error('connection refused')});
|
||||
expect(mapCassandraDriverError(err)).toBe(err);
|
||||
});
|
||||
|
||||
it('returns unrelated errors unchanged', () => {
|
||||
const err = new cassandra.errors.ResponseError(0x2200, 'invalid query');
|
||||
expect(mapCassandraDriverError(err)).toBe(err);
|
||||
});
|
||||
});
|
||||
@@ -1,7 +1,8 @@
|
||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
|
||||
import {ServiceUnavailableError} from '@fluxer/errors/src/domains/core/ServiceUnavailableError';
|
||||
import {getClient} from '@pkgs/cassandra/src/Client';
|
||||
import type cassandra from 'cassandra-driver';
|
||||
import cassandra from 'cassandra-driver';
|
||||
import {Logger} from '../Logger';
|
||||
import {getQueryType, logBatch, logQuery} from './CassandraDevLogger';
|
||||
import {getIsDev} from './CassandraMetaRegistry';
|
||||
@@ -16,6 +17,22 @@ import {
|
||||
|
||||
const DEFAULT_MAX_PARTITION_KEYS_PER_QUERY = 100;
|
||||
|
||||
function isDriverOverloadError(err: unknown): boolean {
|
||||
if (err instanceof cassandra.errors.BusyConnectionError) {
|
||||
return true;
|
||||
}
|
||||
if (err instanceof cassandra.errors.NoHostAvailableError && err.innerErrors) {
|
||||
return Object.values(err.innerErrors as Record<string, unknown>).some(
|
||||
(innerError) => innerError instanceof cassandra.errors.BusyConnectionError,
|
||||
);
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
export function mapCassandraDriverError(err: unknown): unknown {
|
||||
return isDriverOverloadError(err) ? new ServiceUnavailableError() : err;
|
||||
}
|
||||
|
||||
export interface CassandraQueryExecutorForTesting {
|
||||
executeQuery<T = Record<string, unknown>, P extends CassandraParams = CassandraParams>(
|
||||
query: PreparedQuery<P>,
|
||||
@@ -121,7 +138,7 @@ export async function executeQuery<T = Record<string, unknown>, P extends Cassan
|
||||
}
|
||||
const errorMessage = err instanceof Error ? err.message : String(err);
|
||||
Logger.warn({error: errorMessage, query: cql, params: paramSummary}, 'Cassandra query failed');
|
||||
throw err;
|
||||
throw mapCassandraDriverError(err);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -249,10 +266,14 @@ async function executeBatch(queries: Array<BatchQuery>, atomic = true): Promise<
|
||||
counter: false,
|
||||
};
|
||||
const startTime = getIsDev() ? performance.now() : 0;
|
||||
await getClient().batch(
|
||||
queries.map(({query, params}) => ({query, params: normalizeInParams(query, params as CassandraParams)})),
|
||||
options,
|
||||
);
|
||||
try {
|
||||
await getClient().batch(
|
||||
queries.map(({query, params}) => ({query, params: normalizeInParams(query, params as CassandraParams)})),
|
||||
options,
|
||||
);
|
||||
} catch (err: unknown) {
|
||||
throw mapCassandraDriverError(err);
|
||||
}
|
||||
if (getIsDev()) {
|
||||
const durationMs = performance.now() - startTime;
|
||||
logBatch(queries, durationMs);
|
||||
|
||||
Reference in New Issue
Block a user