Monitoring

qianmoQqianmoQ· 更新于 2026-10-08· 阅读 34 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

Metric

Built-in Metrics

We offer data monitoring for built-in metrics, which includes events, actions, and token usage.

Event and Action Metrics

ScopeMetricsDescriptionType
AgentnumOfEventProcessedThe total number of Events this operator has processed.Count
AgentnumOfEventProcessedPerSecThe number of Events this operator has processed per second.Meter
AgentnumOfActionsExecutedThe total number of actions this operator has executed.Count
AgentnumOfActionsExecutedPerSecThe number of actions this operator has executed per second.Meter
Actionaction.<action_name>.numOfActionsExecutedThe total number of actions this operator has executed for a specific action name.Count
Actionaction.<action_name>.numOfActionsExecutedPerSecThe number of actions this operator has executed per second for a specific action name.Meter
AgenteventLogTruncatedEventsNumber of event log records whose payload was truncated at STANDARD level. Increments once per event, regardless of how many fields inside it were truncated. Use this to decide whether to raise truncation thresholds or move specific event types to VERBOSE.Count

Token Usage Metrics

Token usage metrics are automatically recorded when chat models are invoked through ChatModelConnection. These metrics help track LLM API usage and costs.

ScopeMetricsDescriptionType
Modelaction.<action_name>.model.<model_name>.promptTokensThe total number of prompt tokens consumed by the model within an action.Count
Modelaction.<action_name>.model.<model_name>.completionTokensThe total number of completion tokens generated by the model within an action.Count
Modelaction.<action_name>.model.<connection_name>.retryCountThe total number of retries performed for model requests when using ErrorHandlingStrategy.RETRY. Only recorded when at least one retry occurs. See retry-wait-interval.Count
Modelaction.<action_name>.model.<connection_name>.retryWaitSecThe total wait time (in seconds) spent across retries for model requests when using ErrorHandlingStrategy.RETRY.Count

How to add custom metrics

In Flink Agents, users implement their logic by defining custom Actions that respond to various Events throughout the Agent lifecycle. To support user-defined metrics, we introduce two new properties: agent_metric_group and action_metric_group in the RunnerContext. These properties allow users to create or update global metrics and independent metrics for actions. For an introduction to metric types, please refer to the Metric types documentation.

Here is the user case example:

Python

class MyAgent(Agent):
    @action(InputEvent.EVENT_TYPE)
    @staticmethod
    def first_action(event: Event, ctx: RunnerContext):
        start_time = time.time_ns()

        # the action logic
        ...

        # Access the main agent metric group
        metrics = ctx.agent_metric_group

        # Update global metrics
        metrics.get_counter("numInputEvent").inc()
        metrics.get_meter("numInputEventPerSec").mark()

        # Access the per-action metric group
        action_metrics = ctx.action_metric_group
        action_metrics.get_histogram("actionLatencyMs") \
            .update(int(time.time_ns() - start_time) // 1000000)

Java

public class MyAgent extends Agent {

    @Action(listenEventTypes = {InputEvent.EVENT_TYPE})
    public static void firstAction(Event event, RunnerContext ctx) throws Exception {
        InputEvent inputEvent = InputEvent.fromEvent(event);
        long startTime = System.currentTimeMillis();

        // the action logic
        ...

        FlinkAgentsMetricGroup metrics = ctx.getAgentMetricGroup();

        metrics.getCounter("numInputEvent").inc();
        metrics.getMeter("numInputEventPerSec").markEvent();

        FlinkAgentsMetricGroup actionMetrics = ctx.getActionMetricGroup();
        actionMetrics
                .getHistogram("actionLatencyMs")
                .update(System.currentTimeMillis() - startTime);
    }
}

How to check the metrics with Flink executor

Flink agents enable the reporting of metrics to external systems by creating a metric identifier prefix in the format <host>.taskmanager.<tm_id>.<job_name>.<operator_name>.<subtask_index>. Agent-specific metrics use key-value metric groups (e.g., action.<action_name>, model.<model_name>) which are exposed as dimensions/labels in reporters that support them (such as Prometheus). Please refer to Flink Metric Reporters for more details.

Additionally, we can check the metric results in the Flink Job WebUI using the metric identifier prefix <subtask_index>.<operator_name>.

Metric Web UI

Log

The Flink Agents’ log system uses Flink’s logging framework. For more details, please refer to the Flink log system documentation.

How to add log in Flink Agents

For adding logs in Java code, you can refer to Flink documentation. In Python, you can add logs using logging. Here is a specific example:

@action(InputEvent.EVENT_TYPE)
@staticmethod
def process_input(event: Event, ctx: RunnerContext) -> None:
    logging.info("Processing input event: %s", event)
    # the action logic

How to check the logs with Flink executor

We can check the log result in the WebUI of Flink Job:

Log Web UI

Event Log

The system supports two types of event loggers: SLF4J Event Log (default) and File Event Log.

By default, the SLF4J Event Log is used. If baseLogDir is configured, the system automatically switches to the File Event Log.

SLF4J Event Log (Default)

The SLF4J Event Log outputs events through a dedicated SLF4J logger (org.apache.flink.agents.EventLog). On startup, the logger automatically configures log4j2 to write events to a separate file ({log.file}.event-log.log) in Flink’s log directory, making them visible in Flink’s Web UI Logs tab. No manual log4j2 configuration is required.

Because all subtasks on a TaskManager share the same log destination, each record additionally carries jobId, taskName, and subtaskId top-level fields so consumers can still distinguish events from different subtasks. The rest of the record follows the common JSON Format described below.

File Event Log

The File Event Log is a file-based event logging system that stores events in structured files within a flat directory. To use it, configure baseLogDir in your Flink config.yaml.

By default, each event is recorded in JSON Lines (JSONL) format, with one JSON object per line. When prettyPrint is enabled, each event is written as formatted multi-line JSON instead, and the log file is no longer in valid JSONL format.

File Structure

The log files follow a naming convention consistent with Flink’s logging standards and are stored in a flat directory structure:

{baseLogDir}/
├── events-{jobId}-{taskName}-{subtaskId}.log
├── events-{jobId}-{taskName}-{subtaskId}.log
└── events-{jobId}-{taskName}-{subtaskId}.log

By default, all File-based Event Logs are stored in the flink-agents subdirectory under the system temporary directory (java.io.tmpdir). You can override the base log directory with the agent.baseLogDir setting in Flink config.yaml.

JSON Format

The JSON record format described here applies to both the SLF4J and File event loggers. The SLF4J logger adds jobId, taskName, and subtaskId fields on top (see SLF4J Event Log); the File logger encodes those values in the file path instead.

Each record contains a top-level timestamp, the resolved logLevel, and a top-level eventType routing key (mirrors event.eventType), followed by the full event object. The top-level eventType makes it easy for downstream tools (e.g. grep, jq, log shippers) to filter by event type without parsing nested JSON:

{
  "timestamp": "2024-01-15T10:30:00Z",
  "logLevel": "STANDARD",
  "eventType": "_input_event",
  "event": {
    "eventType": "_input_event",
    "...": "..."
  }
}

Event Log Levels

Each event type is logged at a configurable verbosity. Three levels are supported:

LevelBehavior
OFFEvents of this type are not logged.
STANDARDEvents are logged, but the payload may be truncated or summarized to keep logs concise. This is the default.
VERBOSEEvents are logged with the full, untruncated payload.

The global default is set by event-log.level. At STANDARD level, the payload is shrunk along three independent axes — long strings, large arrays, and deep nesting — controlled by event-log.standard.max-string-length, event-log.standard.max-array-elements, and event-log.standard.max-depth respectively. Setting any threshold to 0 disables that specific truncation; setting all three to 0 makes STANDARD behave identically to VERBOSE (apart from the logLevel label). The exact truncation strategy may evolve over time; the contract is only that STANDARD keeps logs concise while VERBOSE preserves the full payload.

Fields that are never truncated. Structural and identifying fields are always preserved in full so log consumers can still group, route, and correlate records: timestamp, logLevel, top-level eventType, and the event’s own id, type, and short scalar fields. Truncation only applies to large nested content (long strings, big arrays, deeply nested objects).

Truncation wrapper format. When a field is truncated at STANDARD level it is replaced by a JSON object that records what was retained and what was dropped. This keeps the record valid JSON and lets downstream tooling detect truncation programmatically:

Truncated contentReplacement
Long string{"truncatedString": "<first N chars>...", "omittedChars": M}
Large array{"truncatedList": [<first N elements>], "omittedElements": M}
Deeply nested object{"truncatedObject": {<scalar fields only>}, "omittedFields": N}

A truncated field changes JSON type (e.g. a string field becomes an object). Consumers that need a stable schema should switch the affected event type to VERBOSE. A counter metric eventLogTruncatedEvents (see Event and Action Metrics) records how often truncation kicked in — use it to decide whether to raise the thresholds or move noisy event types to VERBOSE.

Example record at STANDARD with a long string and a large array truncated:

{
  "timestamp": "2024-01-15T10:30:00Z",
  "logLevel": "STANDARD",
  "eventType": "_chat_request_event",
  "event": {
    "eventType": "_chat_request_event",
    "id": "...",
    "model": "gpt-4",
    "messages": {
      "truncatedList": [
        {"role": "system", "content": "You are a helpful assistant..."},
        {"role": "user", "content": {"truncatedString": "Analyze this doc...", "omittedChars": 1000}}
      ],
      "omittedElements": 30
    }
  }
}

Per-event-type log levels

You can override the level for individual event types using the event-log.type.<EVENT_TYPE>.level config key, where <EVENT_TYPE> is the event’s routing type string (the same string that appears as eventType in the JSON log). Built-in events use short snake-cased names such as:

Event class<EVENT_TYPE> value
InputEvent_input_event
OutputEvent_output_event
ChatRequestEvent_chat_request_event
ChatResponseEvent_chat_response_event
ToolRequestEvent_tool_request_event
ToolResponseEvent_tool_response_event
ContextRetrievalRequestEvent_context_retrieval_request_event
ContextRetrievalResponseEvent_context_retrieval_response_event

Each event type has its own independently overridable key, so a job-level override does not clobber other entries from config.yaml.

Resolution is hierarchical — the resolver walks up dot-separated segments of the event type, mirroring Log4j’s logger hierarchy. For a user-defined event type com.example.myapp.OrderEvent, the lookup order is:

  1. event-log.type.com.example.myapp.OrderEvent.level (exact match)
  2. event-log.type.com.example.myapp.level (package prefix)
  3. event-log.type.com.example.level
  4. … (continues walking up)
  5. event-log.level (global default)
  6. Built-in default: STANDARD

Built-in event type strings like _chat_request_event contain no dots, so they are matched exactly against the configured key.

Example config.yaml:

# Keep all events at STANDARD by default
event-log.level: STANDARD

# Log every ChatRequestEvent with the full payload
event-log.type._chat_request_event.level: VERBOSE

# Skip context retrieval requests entirely
event-log.type._context_retrieval_request_event.level: OFF

Because each event type has its own config key, you can override a single level at job submission time without touching the shared config.yaml:

flink run ... \
  -Devent-log.type._chat_request_event.level=VERBOSE

Other per-type levels from config.yaml are preserved — the -D flag only overrides the one key it names.

Compatibility Notes

  • Default behavior changed. Before this feature, every event was logged in full. The new default is STANDARD, which truncates large payloads. To restore the previous behavior either globally or per type, set the level to VERBOSE.
  • Old log records still parse. Records written before this feature have no logLevel or top-level eventType. They deserialize correctly and are treated as VERBOSE (their payloads were never truncated).

评论

登录后参与评论

正在加载评论…