Configuration

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

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

How to configure Flink Agents

There are two ways to configure Flink Agents, listed in order of priority from high to low:

  1. Setting via the AgentsExecutionEnvironment
  2. Setting via a Flink YAML configuration file

The AgentsExecutionEnvironment applies to Agents from the AgentsExecutionEnvironment, and the Flink YAML configuration file applies to all Flink Agents Jobs using the same configuration file.

In case of duplicate keys, the value from the highest priority will override those from lower priorities.

Setting via the AgentsExecutionEnvironment

Users can explicitly modify the configuration when defining the AgentsExecutionEnvironment:

python

# Get Flink Agents execution environment
agents_env = AgentsExecutionEnvironment.get_execution_environment()

# Get configuration object from the environment
config = agents_env.get_configuration()

# Set custom configuration using a direct key (string-based key)
# This is suitable for user-defined or non-standardized settings.
config.set_int("kafkaActionStateTopicNumPartitions", 128)

# Set framework-level configuration using a predefined ConfigOption class
# This ensures type safety and better integration with the framework.
config.set(AgentExecutionOptions.ERROR_HANDLING_STRATEGY, ErrorHandlingStrategy.RETRY)

java

// Get Flink Agents execution environment
AgentsExecutionEnvironment agentsEnv = AgentsExecutionEnvironment.getExecutionEnvironment(env);

// Get configuration object
Configuration config = agentsEnv.getConfig();

// Set custom configuration using key (direct string key)
config.setInt("kafkaActionStateTopicNumPartitions", 128);  // Kafka topic partitions count

// Set the list of event listeners
config.set(AgentConfigOptions.EVENT_LISTENERS, List.of(MyCustomListener.class.getName()));

// Set framework configuration using ConfigOption (predefined option class)
config.set(AgentExecutionOptions.ERROR_HANDLING_STRATEGY, ErrorHandlingStrategy.RETRY);

Setting via the Flink YAML configuration file

Flink Agents allows reading configurations from the Flink YAML configuration file.

Format

As part of the Flink configuration file, the flink agents configuration must follow this format, with all agent-specific settings nested under the agent key:

agent:
  # Agent-specific configurations
  error-handling-strategy: retry
  chat:
    async: true

Loading Behavior

By default, the configuration is automatically loaded from $FLINK_HOME/conf/config.yaml.

Special Condition

In the following two cases, Flink Agents may not locate the corresponding configuration file, necessitating manual configuration. If the files are not set, no configuration files will be loaded, potentially resulting in unexpected behavior or failures.

  • For MiniCluster: Manual setup is required — always export the environment variable before running the job:

    export FLINK_CONF_DIR="path/to/your/config.yaml"

    This ensures that Flink can locate and load the configuration file correctly.

  • Local mode: When run without flink, use the AgentsExecutionEnvironment.get_configuration() API to load the YAML file directly:

    config = agents_env.get_configuration("path/to/your/config.yaml")

Built-in configuration options

Core Options

Here is the list of all built-in core configuration options.

Key Default Type Description

eventLoggerType SLF4J LoggerType Which built-in event logger to use. Valid values: SLF4J (writes JSON through a dedicated SLF4J logger so events show up in Flink’s Web UI Logs tab) and FILE (writes per-subtask .log files under baseLogDir). Setting baseLogDir overrides this and forces FILE.

baseLogDir (none) String Base directory for file-based event logs. If not set, uses java.io.tmpdir/flink-agents. Setting this value also implicitly switches eventLoggerType to file.

prettyPrint false boolean Whether to enable pretty-printed JSON format for event logs. When set to true, each event is written as formatted multi-line JSON instead of JSONL (JSON Lines) format.

Note: enabling this option makes the log file no longer valid JSONL format.

event-listeners none List<String> The list of event listener class names. Each class must implement the EventListener interface and provide a public no-argument constructor.

Note: Currently, custom event listeners are only supported in Java.

error-handling-strategy ErrorHandlingStrategy.FAIL ErrorHandlingStrategy Strategy for handling errors during model requests, include timeout and unexpected output schema.
The option value could be:

  • ErrorHandlingStrategy.FAIL
  • ErrorHandlingStrategy.RETRY
  • ErrorHandlingStrategy.IGNORE

max-retries 3 int Number of retries when using ErrorHandlingStrategy.RETRY.

retry-wait-interval 1 int Base wait interval in seconds between retries when using ErrorHandlingStrategy.RETRY. Uses exponential backoff: the actual wait time for the Nth retry is retry-wait-interval * 2^(N-1) seconds. For example, with default 1s, waits are 1s, 2s, 4s, etc. Retry count and total wait time are reported in ChatResponseEvent and recorded as metrics (retryCount, retryWaitSec) under the connection name.

chat.async true boolean Whether chat asynchronously for built-in chat action.

tool-call.async true boolean Whether process tool call for built-in tool call action.

rag.async true boolean Whether retrieve context asynchronously for built-in context retrieval action.

num-async-threads os cpu count * 2 int The thread pool size for async executor.

job-identifier none String The unique identifier of job, remaining consistent after restoring from a savepoint. If not set, uses flink job id.

event-log.level STANDARD EventLogLevel Global default verbosity for the Event Log. Valid values: OFF (skip event), STANDARD (payload may be truncated/summarized to keep logs concise), VERBOSE (full payload). Can be overridden per event type — see Per-event-type log levels.

event-log.type.<EVENT_TYPE>.level (inherits) EventLogLevel Override the log level for a specific event type. <EVENT_TYPE> is the event’s routing type string (the same value that appears as eventType in the JSON log, e.g., _chat_request_event for built-ins, or com.example.myapp.OrderEvent for user-defined types). For dotted types, resolution walks up dot segments before falling back to event-log.level. See Per-event-type log levels for examples.

event-log.standard.max-string-length 2000 int At STANDARD level, strings in the event payload longer than this are truncated. Has no effect at VERBOSE.

event-log.standard.max-array-elements 20 int At STANDARD level, arrays in the event payload with more than this many elements are truncated. Has no effect at VERBOSE.

event-log.standard.max-depth 5 int At STANDARD level, objects nested deeper than this are summarized. Has no effect at VERBOSE.

short-term-memory.state-ttl.ms 0 long Time-to-live for short-term memory state in milliseconds. Set to a value greater than 0 to enable TTL; 0 disables it.

short-term-memory.state-ttl.update-type ON_READ_AND_WRITE ShortTermMemoryTtlUpdate Update policy for short-term memory TTL. Only applies when short-term-memory.state-ttl.ms is greater than 0. Valid values: ON_CREATE_AND_WRITE, ON_READ_AND_WRITE.

short-term-memory.state-ttl.visibility NEVER_RETURN_EXPIRED ShortTermMemoryTtlVisibility Visibility policy for expired short-term memory state. Only applies when short-term-memory.state-ttl.ms is greater than 0. Valid values: NEVER_RETURN_EXPIRED, RETURN_EXPIRED_IF_NOT_CLEANED_UP.

Action State Store

Common

KeyDefaultTypeDescription
actionStateStoreBackend(none)StringThe backend for action state store. Supported values: "kafka", "fluss".

Kafka-based Action State Store

Here are the configuration options for Kafka-based Action State Store.

KeyDefaultTypeDescription
kafkaBootstrapServers“localhost:9092”StringThe config parameter specifies the Kafka bootstrap server.
kafkaActionStateTopic(none)StringThe config parameter specifies the Kafka topic for action state.
kafkaActionStateTopicNumPartitions64IntegerThe config parameter specifies the number of partitions for the Kafka action state topic.
kafkaActionStateTopicReplicationFactor1IntegerThe config parameter specifies the replication factor for the Kafka action state topic.

Fluss-based Action State Store

Here are the configuration options for Fluss-based Action State Store.

KeyDefaultTypeDescription
flussBootstrapServers“localhost:9123”StringThe Fluss bootstrap servers address.
flussActionStateDatabase“flink_agents”StringThe Fluss database name for storing action state.
flussActionStateTable(none)StringThe Fluss table name for storing action state.
flussActionStateTableBuckets64IntegerThe number of buckets for the Fluss action state table.
flussSecurityProtocol“PLAINTEXT”StringThe authentication protocol for Fluss client. Valid values: PLAINTEXT (default, no authentication), SASL (SASL/PLAIN authentication).
flussSaslMechanism“PLAIN”StringThe SASL mechanism for Fluss authentication.
flussSaslJaasConfig(none)StringThe JAAS configuration string for Fluss SASL authentication.
flussSaslUsername(none)StringThe username for Fluss SASL authentication.
flussSaslPassword(none)StringThe password for Fluss SASL authentication.

评论

登录后参与评论

正在加载评论…