Skip to content

Commit 8ec5c05

Browse files
committed
feat(otel): cap oversized export batches
1 parent 7b11849 commit 8ec5c05

9 files changed

Lines changed: 445 additions & 15 deletions

File tree

packages/core/src/utils.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ type LangfuseEnvVar =
66
| "LANGFUSE_TIMEOUT"
77
| "LANGFUSE_FLUSH_AT"
88
| "LANGFUSE_FLUSH_INTERVAL"
9+
| "LANGFUSE_OTEL_MAX_BATCH_SIZE_BYTES"
910
| "LANGFUSE_MEDIA_UPLOAD_ENABLED"
1011
| "LANGFUSE_LOG_LEVEL"
1112
| "LANGFUSE_DEBUG"

packages/otel/README.md

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,8 +19,16 @@ npm install @langfuse/otel @opentelemetry/sdk-trace-node
1919
LANGFUSE_PUBLIC_KEY="pk-lf-..."
2020
LANGFUSE_SECRET_KEY="sk-lf-..."
2121
LANGFUSE_BASE_URL="https://cloud.langfuse.com" # 🇪🇺 EU region. 🇺🇸 US: https://us.cloud.langfuse.com
22+
LANGFUSE_OTEL_MAX_BATCH_SIZE_BYTES="67108864" # Optional; defaults to 64 MiB
2223
```
2324

25+
The batch byte limit applies to the final serialized OTLP/HTTP JSON request
26+
before compression. If a batch exceeds the limit, the entire batch is dropped;
27+
it is not split, retried, or truncated. Invalid values fall back to 64 MiB with
28+
a warning. Set the limit to a positive decimal safe integer; surrounding
29+
whitespace is ignored. Custom exporters passed to `LangfuseSpanProcessor`
30+
bypass this limit because they may use a different wire format or transport.
31+
2432
## Quickstart
2533

2634
```typescript

packages/otel/package.json

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -41,7 +41,8 @@
4141
"dist"
4242
],
4343
"dependencies": {
44-
"@langfuse/core": "workspace:^"
44+
"@langfuse/core": "workspace:^",
45+
"@opentelemetry/otlp-transformer": "0.220.0"
4546
},
4647
"peerDependencies": {
4748
"@opentelemetry/api": "^1.9.0",
Lines changed: 123 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,123 @@
1+
import { getEnv, getGlobalLogger } from "@langfuse/core";
2+
import { ExportResultCode } from "@opentelemetry/core";
3+
import { JsonTraceSerializer } from "@opentelemetry/otlp-transformer";
4+
import type { ReadableSpan, SpanExporter } from "@opentelemetry/sdk-trace-base";
5+
6+
export const DEFAULT_MAX_BATCH_SIZE_BYTES = 64 * 1024 * 1024;
7+
8+
export function getSerializedBatchSizeBytes(
9+
spans: ReadableSpan[],
10+
): number | undefined {
11+
return JsonTraceSerializer.serializeRequest(spans)?.byteLength;
12+
}
13+
14+
export function resolveMaxBatchSizeBytes(rawValue: string | undefined): number {
15+
if (rawValue === undefined) return DEFAULT_MAX_BATCH_SIZE_BYTES;
16+
17+
const normalizedValue = rawValue.trim();
18+
if (!/^\d+$/.test(normalizedValue)) {
19+
getGlobalLogger().warn(
20+
"Invalid LANGFUSE_OTEL_MAX_BATCH_SIZE_BYTES. Using the default limit.",
21+
{ defaultMaxBatchSizeBytes: DEFAULT_MAX_BATCH_SIZE_BYTES },
22+
);
23+
return DEFAULT_MAX_BATCH_SIZE_BYTES;
24+
}
25+
26+
const parsedValue = Number(normalizedValue);
27+
if (!Number.isSafeInteger(parsedValue) || parsedValue <= 0) {
28+
getGlobalLogger().warn(
29+
"Invalid LANGFUSE_OTEL_MAX_BATCH_SIZE_BYTES. Using the default limit.",
30+
{ defaultMaxBatchSizeBytes: DEFAULT_MAX_BATCH_SIZE_BYTES },
31+
);
32+
return DEFAULT_MAX_BATCH_SIZE_BYTES;
33+
}
34+
35+
return parsedValue;
36+
}
37+
38+
export function resolveMaxBatchSizeBytesFromEnvironment(): number {
39+
const processValue =
40+
typeof process !== "undefined"
41+
? process.env.LANGFUSE_OTEL_MAX_BATCH_SIZE_BYTES
42+
: undefined;
43+
44+
return resolveMaxBatchSizeBytes(
45+
processValue !== undefined
46+
? processValue
47+
: getEnv("LANGFUSE_OTEL_MAX_BATCH_SIZE_BYTES"),
48+
);
49+
}
50+
51+
export class SizeLimitedSpanExporter implements SpanExporter {
52+
private readonly delegate: SpanExporter;
53+
private readonly maxBatchSizeBytes: number;
54+
private readonly serializeRequest: (
55+
spans: ReadableSpan[],
56+
) => { byteLength: number } | undefined;
57+
58+
constructor(params: {
59+
delegate: SpanExporter;
60+
maxBatchSizeBytes: number;
61+
serializeRequest?: (
62+
spans: ReadableSpan[],
63+
) => { byteLength: number } | undefined;
64+
}) {
65+
this.delegate = params.delegate;
66+
this.maxBatchSizeBytes = params.maxBatchSizeBytes;
67+
this.serializeRequest =
68+
params.serializeRequest ?? JsonTraceSerializer.serializeRequest;
69+
}
70+
71+
export(
72+
spans: ReadableSpan[],
73+
resultCallback: Parameters<SpanExporter["export"]>[1],
74+
): void {
75+
let serializedSizeBytes: number | undefined;
76+
77+
try {
78+
serializedSizeBytes = this.serializeRequest(spans)?.byteLength;
79+
} catch (error) {
80+
resultCallback({
81+
code: ExportResultCode.FAILED,
82+
error: error instanceof Error ? error : new Error(String(error)),
83+
});
84+
return;
85+
}
86+
87+
if (serializedSizeBytes === undefined) {
88+
resultCallback({
89+
code: ExportResultCode.FAILED,
90+
error: new Error("Failed to serialize OpenTelemetry span batch."),
91+
});
92+
return;
93+
}
94+
95+
if (serializedSizeBytes > this.maxBatchSizeBytes) {
96+
getGlobalLogger().warn(
97+
"Dropping OpenTelemetry span batch because its serialized request exceeds the configured byte limit.",
98+
{
99+
maxBatchSizeBytes: this.maxBatchSizeBytes,
100+
serializedSizeBytes,
101+
spanCount: spans.length,
102+
},
103+
);
104+
resultCallback({
105+
code: ExportResultCode.FAILED,
106+
error: new Error(
107+
`Serialized OpenTelemetry span batch size ${serializedSizeBytes} bytes exceeds the configured limit of ${this.maxBatchSizeBytes} bytes.`,
108+
),
109+
});
110+
return;
111+
}
112+
113+
this.delegate.export(spans, resultCallback);
114+
}
115+
116+
forceFlush(): Promise<void> {
117+
return this.delegate.forceFlush?.() ?? Promise.resolve();
118+
}
119+
120+
shutdown(): Promise<void> {
121+
return this.delegate.shutdown();
122+
}
123+
}

packages/otel/src/span-processor.ts

Lines changed: 22 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,10 @@ import {
2323
} from "@opentelemetry/sdk-trace-base";
2424

2525
import { MediaService } from "./MediaService.js";
26+
import {
27+
SizeLimitedSpanExporter,
28+
resolveMaxBatchSizeBytesFromEnvironment,
29+
} from "./size-limited-span-exporter.js";
2630
import { isDefaultExportSpan } from "./span-filter.js";
2731

2832
/**
@@ -78,6 +82,8 @@ export type ShouldExportSpan = (params: { otelSpan: ReadableSpan }) => boolean;
7882
export interface LangfuseSpanProcessorParams {
7983
/**
8084
* Custom OpenTelemetry span exporter. If not provided, a default OTLP exporter will be used.
85+
* Custom exporters bypass Langfuse's OTLP batch byte limit and must enforce
86+
* any transport-specific request limits themselves.
8187
*/
8288
exporter?: SpanExporter;
8389

@@ -275,19 +281,22 @@ export class LangfuseSpanProcessor implements SpanProcessor {
275281
? !["false", "0"].includes(envMediaUploadEnabled.toLowerCase())
276282
: true);
277283

278-
const exporter =
279-
params?.exporter ??
280-
new OTLPTraceExporter({
281-
url: `${baseUrl}/api/public/otel/v1/traces`,
282-
headers: {
283-
Authorization: `Basic ${authHeaderValue}`,
284-
"x-langfuse-sdk-name": "javascript",
285-
"x-langfuse-sdk-version": LANGFUSE_SDK_VERSION,
286-
"x-langfuse-public-key": publicKey ?? "<missing>",
287-
...params?.additionalHeaders,
288-
},
289-
timeoutMillis: timeoutSeconds * 1_000,
290-
});
284+
const exporter = params?.exporter
285+
? params.exporter
286+
: new SizeLimitedSpanExporter({
287+
maxBatchSizeBytes: resolveMaxBatchSizeBytesFromEnvironment(),
288+
delegate: new OTLPTraceExporter({
289+
url: `${baseUrl}/api/public/otel/v1/traces`,
290+
headers: {
291+
Authorization: `Basic ${authHeaderValue}`,
292+
"x-langfuse-sdk-name": "javascript",
293+
"x-langfuse-sdk-version": LANGFUSE_SDK_VERSION,
294+
"x-langfuse-public-key": publicKey ?? "<missing>",
295+
...params?.additionalHeaders,
296+
},
297+
timeoutMillis: timeoutSeconds * 1_000,
298+
}),
299+
});
291300

292301
this.processor =
293302
params?.exportMode === "immediate"

pnpm-lock.yaml

Lines changed: 3 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

tests/integration/span-processor.integration.test.ts

Lines changed: 57 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,14 @@
1+
import { createServer } from "node:http";
2+
import type { AddressInfo } from "node:net";
3+
14
import {
25
LogLevel,
36
configureGlobalLogger,
47
resetGlobalLogger,
58
} from "@langfuse/core";
9+
import { LangfuseSpanProcessor } from "@langfuse/otel";
610
import { trace } from "@opentelemetry/api";
11+
import { NodeSDK } from "@opentelemetry/sdk-node";
712
import { startObservation } from "@langfuse/tracing";
813
import { describe, it, expect, beforeEach, afterEach, vi } from "vitest";
914

@@ -28,6 +33,7 @@ describe("LangfuseSpanProcessor E2E Tests", () => {
2833

2934
afterEach(async () => {
3035
delete process.env.LANGFUSE_MEDIA_UPLOAD_ENABLED;
36+
delete process.env.LANGFUSE_OTEL_MAX_BATCH_SIZE_BYTES;
3137
await teardownTestEnvironment(testEnv);
3238
vi.restoreAllMocks();
3339
resetGlobalLogger();
@@ -718,6 +724,57 @@ describe("LangfuseSpanProcessor E2E Tests", () => {
718724
});
719725

720726
describe("Export Mode Selection", () => {
727+
it("does not send an oversized request through the default exporter", async () => {
728+
await teardownTestEnvironment(testEnv);
729+
process.env.LANGFUSE_OTEL_MAX_BATCH_SIZE_BYTES = "1";
730+
const warn = vi
731+
.spyOn(console, "warn")
732+
.mockImplementation(() => undefined);
733+
let requestCount = 0;
734+
const server = createServer((request, response) => {
735+
requestCount += 1;
736+
request.resume();
737+
response.writeHead(200);
738+
response.end();
739+
});
740+
await new Promise<void>((resolve) =>
741+
server.listen(0, "127.0.0.1", resolve),
742+
);
743+
const { port } = server.address() as AddressInfo;
744+
const spanProcessor = new LangfuseSpanProcessor({
745+
baseUrl: `http://127.0.0.1:${port}`,
746+
publicKey: "pk-lf-test",
747+
secretKey: "sk-lf-test",
748+
exportMode: "immediate",
749+
mediaUploadEnabled: false,
750+
});
751+
const sdk = new NodeSDK({
752+
spanProcessor,
753+
instrumentations: [],
754+
});
755+
sdk.start();
756+
757+
try {
758+
const span = startObservation("oversized-default-export");
759+
span.end();
760+
await spanProcessor.forceFlush();
761+
762+
expect(requestCount).toBe(0);
763+
expect(warn).toHaveBeenCalledWith(
764+
expect.stringContaining("Dropping OpenTelemetry span batch"),
765+
expect.objectContaining({
766+
maxBatchSizeBytes: 1,
767+
spanCount: 1,
768+
}),
769+
);
770+
} finally {
771+
await sdk.shutdown();
772+
await new Promise<void>((resolve, reject) =>
773+
server.close((error) => (error ? reject(error) : resolve())),
774+
);
775+
}
776+
});
777+
721778
it("should use BatchSpanProcessor by default", async () => {
722779
// Default testEnv uses batched mode
723780
const span1 = startObservation("default-batch-1");

0 commit comments

Comments
 (0)