Fixed batched otel in dev and made prod configurable to be batched as well

This commit is contained in:
Eric Allam
2024-04-02 12:58:56 +01:00
parent d876c358d2
commit eb60126284
5 changed files with 60 additions and 7 deletions
+5
View File
@@ -0,0 +1,5 @@
---
"@trigger.dev/core": patch
---
Fixed batch otel flushing
+11
View File
@@ -118,6 +118,17 @@ const EnvironmentSchema = z.object({
DEV_OTEL_LOG_SCHEDULED_DELAY_MILLIS: z.string().default("200"),
DEV_OTEL_LOG_EXPORT_TIMEOUT_MILLIS: z.string().default("30000"),
DEV_OTEL_LOG_MAX_QUEUE_SIZE: z.string().default("512"),
PROD_OTEL_BATCH_PROCESSING_ENABLED: z.string().default("0"),
PROD_OTEL_SPAN_MAX_EXPORT_BATCH_SIZE: z.string().default("64"),
PROD_OTEL_SPAN_SCHEDULED_DELAY_MILLIS: z.string().default("200"),
PROD_OTEL_SPAN_EXPORT_TIMEOUT_MILLIS: z.string().default("30000"),
PROD_OTEL_SPAN_MAX_QUEUE_SIZE: z.string().default("512"),
PROD_OTEL_LOG_MAX_EXPORT_BATCH_SIZE: z.string().default("64"),
PROD_OTEL_LOG_SCHEDULED_DELAY_MILLIS: z.string().default("200"),
PROD_OTEL_LOG_EXPORT_TIMEOUT_MILLIS: z.string().default("30000"),
PROD_OTEL_LOG_MAX_QUEUE_SIZE: z.string().default("512"),
RUNTIME_WAIT_THRESHOLD_IN_MS: z.coerce.number().int().default(30000),
// Internal OTEL environment variables
@@ -477,6 +477,46 @@ export class EnvironmentVariablesRepository implements Repository {
key: "TRIGGER_RUNTIME_WAIT_THRESHOLD_IN_MS",
value: String(env.RUNTIME_WAIT_THRESHOLD_IN_MS),
},
...(env.PROD_OTEL_BATCH_PROCESSING_ENABLED === "1"
? [
{
key: "OTEL_BATCH_PROCESSING_ENABLED",
value: "1",
},
{
key: "OTEL_SPAN_MAX_EXPORT_BATCH_SIZE",
value: env.PROD_OTEL_SPAN_MAX_EXPORT_BATCH_SIZE,
},
{
key: "OTEL_SPAN_SCHEDULED_DELAY_MILLIS",
value: env.PROD_OTEL_SPAN_SCHEDULED_DELAY_MILLIS,
},
{
key: "OTEL_SPAN_EXPORT_TIMEOUT_MILLIS",
value: env.PROD_OTEL_SPAN_EXPORT_TIMEOUT_MILLIS,
},
{
key: "OTEL_SPAN_MAX_QUEUE_SIZE",
value: env.PROD_OTEL_SPAN_MAX_QUEUE_SIZE,
},
{
key: "OTEL_LOG_MAX_EXPORT_BATCH_SIZE",
value: env.PROD_OTEL_LOG_MAX_EXPORT_BATCH_SIZE,
},
{
key: "OTEL_LOG_SCHEDULED_DELAY_MILLIS",
value: env.PROD_OTEL_LOG_SCHEDULED_DELAY_MILLIS,
},
{
key: "OTEL_LOG_EXPORT_TIMEOUT_MILLIS",
value: env.PROD_OTEL_LOG_EXPORT_TIMEOUT_MILLIS,
},
{
key: "OTEL_LOG_MAX_QUEUE_SIZE",
value: env.PROD_OTEL_LOG_MAX_QUEUE_SIZE,
},
]
: []),
];
}
@@ -92,10 +92,7 @@ export class CompleteAttemptService extends BaseService {
},
});
logger.debug("Completed attempt successfully, ACKing message", {
serializedOutput: completion.output,
outputType: completion.outputType,
});
logger.debug("Completed attempt successfully, ACKing message");
await marqs?.acknowledgeMessage(taskRunAttempt.taskRunId);
+3 -3
View File
@@ -91,7 +91,7 @@ export class TracingSDK {
public readonly asyncResourceDetector = new AsyncResourceDetector();
private readonly _logProvider: LoggerProvider;
private readonly _spanExporter: SpanExporter;
private readonly _traceProvider: TracerProvider;
private readonly _traceProvider: NodeTracerProvider;
public readonly getLogger: LoggerProvider["getLogger"];
public readonly getTracer: TracerProvider["getTracer"];
@@ -195,12 +195,12 @@ export class TracingSDK {
}
public async flush() {
await this._spanExporter.forceFlush?.();
await this._traceProvider.forceFlush();
await this._logProvider.forceFlush();
}
public async shutdown() {
await this._spanExporter.shutdown();
await this._traceProvider.shutdown();
await this._logProvider.shutdown();
}
}