MCP
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.
- In Python, use
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 tokensBasicAuth- For username/password authenticationApiKeyAuth- 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.
| Argument | Default | Description |
|---|---|---|
maxRetries | 3 | Maximum number of retries for a failed remote call. |
initialBackoffMs | 100 | Initial backoff before the first retry, in milliseconds. |
maxBackoffMs | 10000 | Upper 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 Language | JDK Version | Default Implementation |
|---|---|---|
| Python | Any | Python SDK |
| Java | JDK 17+ | Java SDK |
| Java | JDK 16 and below | Python 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() {
// ...
}评论
登录后参与评论
KnowForge