[integration][observability] Add OpenTelemetry GenAI exporter for Agent Traces (Phase 1) - #1160
Zhuoxi2000 wants to merge 7 commits into
Conversation
sangkyoonnam
left a comment
There was a problem hiding this comment.
Thanks for taking the suggestions in. I pulled the branch, built the module and ran its tests (18 + 4 pass), and checked the span shapes against semantic-conventions-genai at e57c543b: the root invoke_agent INTERNAL span, chat as CLIENT, execute_tool as INTERNAL with gen_ai.tool.call.id and gen_ai.tool.type, and error.type carrying the recorded error type (which the built-in reporters derive from the root-cause class) all match the internal-agent and tool span tables. gen_ai.provider.name is Required on the chat span and stays unset because the record doesn't carry it; the docs say so, and that's the one gap I'd keep stated rather than fixed here. None of this blocks. Three things worth a look below (the usage mapping, request size, and the conversation id), plus two nits.
| private ExportSummary export(List<TraceRecord> records, List<ConverterDiagnostic> diagnostics) { | ||
| List<SpanData> spans = assembler.assemble(records, diagnostics); | ||
| if (!spans.isEmpty()) { | ||
| CompletableResultCode result = spanExporter.export(spans); |
There was a problem hiding this comment.
This exports every assembled span in one call. I checked the pinned OpenTelemetry 1.51.0 gRPC and HTTP exporters: neither splits the collection, each marshals it into a single request. The OTel Collector's default gRPC receive limit is 4 MiB (grpc-go's default unless max_recv_msg_size_mib is set), so an oversized export is rejected as a whole. Could exports use bounded batches and report partial progress if a later batch fails? The SDK's 512-span BatchSpanProcessor default is one reference point, though a span count alone doesn't guarantee a request stays under the byte limit. If batching is added, should the 30 s wait at L170 apply per batch or to the whole run?
There was a problem hiding this comment.
Done in 36e9717: 512-span batches, 30 s per batch, and failures report delivered spans.
| } | ||
|
|
||
| private static long epochNanos(String isoTimestamp) { | ||
| Instant instant = Instant.parse(isoTimestamp); |
There was a problem hiding this comment.
nit: an otherwise valid lifecycle record with an invalid timestamp throws DateTimeParseException through RunAccumulator.accept and aborts the whole export, while a JSON decoding failure becomes a MALFORMED_RECORD diagnostic. EventContext writes Instant.now().toString() on the built-in path, but the other constructor and the timestamped report methods take the string as given. Could the timestamp be validated before the record enters either accumulator, with a diagnostic and skip instead? A test with one invalid timestamp followed by a valid record would cover it.
There was a problem hiding this comment.
Done: an invalid timestamp now yields a MALFORMED_RECORD diagnostic and skips that record.
| if (eventAttributes == null) { | ||
| return; | ||
| } | ||
| Object prompt = eventAttributes.get("promptTokens"); |
There was a problem hiding this comment.
executionFinished() produces no usage attributes: the Java and Python chat models keep token counts in the response's extraArgs and in metric counters (BaseChatModelSetup.java:181), and nothing writes promptTokens / completionTokens into a terminal record's eventAttributes. The branches here are exercised only by the test fixtures. The docs already defer richer usage metrics, so this is a clarification rather than a request to grow the PR: could the mapping docs say that logs from the current built-in reporters don't supply the terminal attributes this converter needs for gen_ai.usage.*, so "when recorded" doesn't read as "usually"?
There was a problem hiding this comment.
Good catch; the docs now say built-in reporters don't write these token attributes yet.
| // the framework running the tool, whatever transport the tool itself uses. | ||
| kind = SpanKind.INTERNAL; | ||
| attributes.put(GEN_AI_OPERATION_NAME, "execute_tool"); | ||
| attributes.put(GEN_AI_TOOL_NAME, entityName); |
There was a problem hiding this comment.
nit: the execute_tool table lists gen_ai.agent.name as conditionally required when applicable, and TraceRecord exposes agentName, but it's only set on the root span. Could tool spans carry it when the record has one? That keeps a tool span self-describing when a backend shows it without its root.
There was a problem hiding this comment.
Done: tool spans now carry gen_ai.agent.name whenever their records have one.
| } | ||
| attributes.put(FA_ENTITY_NAME, entityName); | ||
| if (any.getBusinessKey() != null) { | ||
| attributes.put(GEN_AI_CONVERSATION_ID, any.getBusinessKey()); |
There was a problem hiding this comment.
businessKey is the keyed-stream key rendered as text. It is a conversation in a chat pipeline and an order or device id elsewhere, while the convention asks for a conversation identifier the library actually has. Would you carry it as flink_agents.business_key unconditionally and map it to gen_ai.conversation.id behind an option? Same for the root mapping at L163.
There was a problem hiding this comment.
Agreed; it's always flink_agents.business_key now, and conversation id is opt-in.
…nt Traces (Phase 1) See apache#970.
…om the current Event Log format
… the opentelemetry-bom
…lation attributes
…s and unwrap truncated errors
…mps, opt-in conversation id, tool agent names
224b58d to
36e9717
Compare
|
All five are addressed in 36e9717. I rebuilt the module at that commit and its 28 tests pass. Exports now go out in bounded batches and stop at the first failed or timed-out one, reporting the spans from the batches that succeeded before it, which answers my batching point. Nothing more from my side. |
|
Thanks for the careful review; @wenjin272, could you take a look when you have time? |
|
Thanks for the PR. Since @joeyutong designed most of Flink Agents’ observability, I think he would be best placed to review this. |
Linked issue: #970
Purpose of change
An Agent Trace Event Log (
event-log.trace.enabled: true) from a Java or Python agent can now be exported to any OTLP backend as OpenTelemetry GenAI traces:This is a new optional module,
flink-agents-integrations-observability-otel, and is not bundled intodist. The goal is to make the execution hierarchy already recorded in Agent Trace visible in standard tracing tools without adding runtime overhead.Runtime flow
exportFilesexpands directories to theirevents-*.logfiles and streams JSON records intoTraceRecord(unknown fields are ignored).AgentTraceSpans.assemblekeeps lifecycle records withexecutionId,inputRunId, and a parseabletimestamp, grouped by execution and input run.invoke_agentroot. Each execution becomes a span under itsparentExecutionId, or under the root if no parent is present, based onentityType.SpanExporterin batches of at most 512 (--batch-size). Diagnostics are logged and returned inExportSummary.Key decisions
inputRunId/executionId): re-exporting the same records produces the same IDs; delivery remains at-least-once.gen_ai.request.modelcomes fromentityMetadata.model(entityNameis the ChatModel resource).gen_ai.provider.nameis left unset because the record does not contain it.execute_toolspans remain INTERNAL regardless of the tool transport.businessKeyis the keyed-stream key, a conversation only in chat pipelines, so it is alwaysflink_agents.business_keyand becomesgen_ai.conversation.idonly when enabled.Behavioral Semantics
Interaction decisions
startedis used as the start record when present, otherwisecreated;createdis never treated as a terminal record.finishedfailederror.typecreatedonly)incomplete=trueINCOMPLETE_EXECUTIONreusedincomplete=trueMISSING_STARTreusedonlystatus=reusedBehavioral contracts
inputRunId, rooted atinvoke_agent {agentName}(INTERNAL), withgen_ai.agent.nameand, when present,flink_agents.business_key;gen_ai.conversation.id=businessKeyonly with--business-key-as-conversation-id.parentExecutionIdspan; otherwise it is attached to the run root.llm→chat {model}(CLIENT):gen_ai.request.modelfromentityMetadata.model, andgen_ai.usage.*_tokensfrompromptTokens/completionTokenswhen the terminal record carries them (the built-in reporters do not yet). Without a model, the span is namedchatand no request model is set.tool→execute_tool {name}(INTERNAL):gen_ai.tool.name;gen_ai.tool.call.id=externalId, otherwisetoolCallId;gen_ai.tool.type=function,extension(remote_function,mcp), or unset (model_built_in). The raw value is kept inflink_agents.tool.type.gen_ai.agent.nameis set when the tool's records carryagentName.action→action {name};parser→parse {name}withgen_ai.operation.name=parse; both are INTERNAL.error.type= the recorded error type, falling back toproblemCategory.flink_agents.{input_run_id, execution_id, entity_type, entity_name, execution.status}, plusflink_agents.business_keywhen present.executionId,inputRunId, ortimestamp, are skipped; a record whose timestamp does not parse is skipped with aMALFORMED_RECORDdiagnostic.Failure behavior
MALFORMED_RECORDdiagnostic naming the file. Reading continues, but a broken stream or 1,000 malformed records stop that file.ExportFailedExceptionwith the number of spans the earlier batches delivered; there is no retry, and re-running is safe.--protocolthrowsIllegalArgumentExceptionat build time.--batch-size, or no input, prints usage and exits 2; a missing input path throwsNoSuchFileException.Tests
AgentTraceSpansTest.testRunBecomesSingleTraceWithRootSpan,testBusinessKeyAsConversationIdIsOptIntestParenting,testDiagnosticsForIncompleteAndMissingStarttestGenAiAttributes,testChatWithoutModeltestGenAiAttributes,testToolTypeMapping,testToolSpanCarriesAgentNametestParserMapping,testParentingtestGenAiAttributes,testErrorTypeFallsBackToProblemCategory,testTruncatedErrorMessagetestCorrelationAttributestestDeterministicIds,testRecordOrderDoesNotMattertestNonLifecycleRecordsIgnored,testInvalidTimestampIsSkippedWithDiagnostic,EventLogOTelExporterTest.testExportJsonlFiletestCreatedStartedTerminal,testCreatedOnly,testCreatedThenFailedWithoutStart,testIncompleteExecution,testDiagnosticsForIncompleteAndMissingStart,testReusedExecutiontestMalformedRecordDiagnostic,testUnsupportedProtocol,testDirectoryDiscoverytestExportsInBoundedBatches,testFailedBatchReportsPartialProgress,testRejectsNonPositiveBatchSizeSpan mapping is covered for each entity type and lifecycle case; input handling uses real files.
Not verified:
SpanExporter, so transport, the receiver's size limit, and the 30 s per-batch timeout are not covered (a failing batch is).Implementation invariants and supporting evidence
EventLogRecordJsonSerializer:entityMetadatais an object,status/problemCategoryare top-level fields, anderrorType/errorMessageare insideeventAttributes. A failed-tool record and an LLM record produced by the real serializer (EventLogRecord+ExecutionLifecycleEvents) assemble into the same spans as the fixtures.OTelIds: trace id = first 16 bytes of SHA-256(inputRunId), span id = first 8 bytes of SHA-256(executionId); the run root hashesinputRunIdwith a domain suffix so it cannot collide with an execution span id; an all-zero id is adjusted to remain valid.opentelemetry-bom(all resolve to 1.51.0).spotless:check, andverifyfor the module and its upstream modules.API
No existing API, record format, or runtime behavior changes.
distis unchanged, so nothing is added to a job's classpath.New public entry points:
EventLogOTelExporter(builder andmain),AgentTraceSpans, andTraceRecord.Existing trace-enabled Event Logs are consumed as-is. The GenAI conventions are still at development stability, so the exported attribute set is pinned per release.
Documentation
doc-neededdoc-not-neededdoc-includedWas this patch authored or co-authored using generative AI tooling?
If yes, include a
Generated-by: <tool name and version> (<model name and version>)line, for exampleGenerated-by: Claude Code 2.1.226 (Claude Opus 4.6), in the commit message so it reaches Git history. Repeat the same line here for reviewer visibility. See the ASF generative tooling guidance.Generated-by: Claude Code 2.1.259 (Claude Fable 5, Claude Opus 5.5)