MCP

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

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

Overview

MCP (Model Context Protocol) is a standardized protocol for integrating AI applications with external data sources and tools. Flink Agents provides the support for using prompts and tools from MCP server.

Declare MCP Server in Agent

Developer can declare a mcp server by decorator/annotation when creating an Agent.

Python

class ReviewAnalysisAgent(Agent):

    @mcp_server
    @staticmethod
    def my_mcp_server() -> ResourceDescriptor:
        """Define MCP server connection."""
        return ResourceDescriptor(clazz=ResourceName.MCP_SERVER,
                                  endpoint="http://127.0.0.1:8000/mcp")

Java

public class ReviewAnalysisAgent extends Agent {

    @MCPServer
    public static ResourceDescriptor myMcp() {
        return ResourceDescriptor.Builder.newBuilder(ResourceName.MCP_SERVER)
                    .addInitialArgument("endpoint", MCP_ENDPOINT)
                    .addInitialArgument("timeout", 30)
                    .build();
    }
}

Key points:

  • Use decorator/annotation to define MCP server connection

    • In Python, use @mcp_server.
    • In Java, use @MCPServer.
  • Use the builder pattern in Java to configure the MCP server with endpoint, timeout, headers, and authentication

Authentication

MCP servers can be configured with authentication:

Python

@mcp_server
@staticmethod
def authenticated_mcp_server() -> ResourceDescriptor:
    """Connect to MCP server with authentication."""
    return ResourceDescriptor(clazz=ResourceName.MCP_SERVER,
                              endpoint="http://api.example.com/mcp",
                              headers={"Authorization": "Bearer your-token"})
    # Or using Basic Authentication
    # credentials = base64.b64encode(b"username:password").decode("ascii")
    # headers={"Authorization": f"Basic {credentials}"}

    # Or using API Key Authentication
    # headers={"X-API-Key": "your-api-key"}

Java

@MCPServer
public static ResourceDescriptor authenticatedMcpServer() {
    // Using Bearer Token Authentication
    return ResourceDescriptor.Builder.newBuilder(ResourceName.MCP_SERVER)
                    .addInitialArgument("endpoint", "http://api.example.com/mcp")
                    .addInitialArgument("timeout", 30)
                    .addInitialArgument("auth", new BearerTokenAuth("your-oauth-token"))
                    .build();

    // Or using Basic Authentication
    // .addInitialArgument("auth", new BasicAuth("username", "password"))

    // Or using API Key Authentication
    // .addInitialArgument("auth", new ApiKeyAuth("X-API-Key", "your-api-key"))
}

Authentication options:

  • BearerTokenAuth - For OAuth 2.0 and JWT tokens
  • BasicAuth - For username/password authentication
  • ApiKeyAuth - For API key authentication via custom headers

Retries

Remote calls to an MCP server (listing/calling tools and prompts) are retried with exponential backoff on transient failures. You can tune the retry behavior per server with the following descriptor arguments.

Note: retry is currently only supported by the Java SDK.

ArgumentDefaultDescription
maxRetries3Maximum number of retries for a failed remote call.
initialBackoffMs100Initial backoff before the first retry, in milliseconds.
maxBackoffMs10000Upper bound for the backoff between retries, in milliseconds.
@MCPServer
public static ResourceDescriptor myMcp() {
    return ResourceDescriptor.Builder.newBuilder(ResourceName.MCP_SERVER)
                .addInitialArgument("endpoint", MCP_ENDPOINT)
                .addInitialArgument("timeout", 30)
                .addInitialArgument("maxRetries", 5)
                .addInitialArgument("initialBackoffMs", 200)
                .addInitialArgument("maxBackoffMs", 10000)
                .build();
}

Use MCP prompts and tools in Agent

MCP prompts and tools are managed by external MCP servers and automatically discovered when you define an MCP server connection in your agent.

Python

class ReviewAnalysisAgent(Agent):

    @mcp_server
    @staticmethod
    def review_mcp_server() -> ResourceDescriptor:
        """Connect to MCP server."""
        return ResourceDescriptor(clazz=ResourceName.MCP_SERVER,
                                  endpoint="http://127.0.0.1:8000/mcp")

    @chat_model_setup
    @staticmethod
    def review_model() -> ResourceDescriptor:
        return ResourceDescriptor(
            clazz=ResourceName.ChatModel.OLLAMA_SETUP,
            connection="ollama_server",
            model="qwen3:8b",
            # Reference MCP prompt by name like local prompt
            prompt="review_analysis_prompt",
            # Reference MCP tool by name like function tool
            tools=["notify_shipping_manager"],
        )

Java

public class ReviewAnalysisAgent extends Agent {

    @MCPServer
    public static ResourceDescriptor myMcp() {
        return ResourceDescriptor.Builder.newBuilder(ResourceName.MCP_SERVER)
                    .addInitialArgument("endpoint", "http://127.0.0.1:8000/mcp")
                    .addInitialArgument("timeout", 30)
                    .build();
    }

    @ChatModelSetup
    public static ResourceDescriptor reviewModel() {
        return ResourceDescriptor.Builder.newBuilder(ResourceName.ChatModel.OLLAMA_SETUP)
                .addInitialArgument("connection", "ollamaChatModelConnection")
                .addInitialArgument("model", "qwen3:8b")
                // Reference MCP prompt by name like local prompt
                .addInitialArgument("prompt", "review_analysis_prompt")
                // Reference MCP tool by name like function tool
                .addInitialArgument("tools", Collections.singletonList("notifyShippingManager"))
                .build();
    }
}

Key points:

  • All tools and prompts from the MCP server are automatically registered.
  • Reference MCP prompts and tools by their names, like reference local prompt and function tool .

Appendix

MCP SDK

Flink Agents offers two implementations of MCP support, based on MCP SDKs in different languages (Python and Java). Typically, users do not need to be aware of this, as the framework automatically determines the appropriate implementation based on the language and version. The default behavior is described as follows:

Agent LanguageJDK VersionDefault Implementation
PythonAnyPython SDK
JavaJDK 17+Java SDK
JavaJDK 16 and belowPython SDK

As shown in the table above, for Java agents running on JDK 17+, the framework automatically uses the Java SDK implementation. If you need to use the Python SDK instead (not recommended), you can set the lang parameter to "python" in the @MCPServer annotation:

@MCPServer(lang = "python")
public static ResourceDescriptor myMcp() {
    // ...
}

评论

登录后参与评论

正在加载评论…