云计算百科
云计算领域专业知识百科平台

【AI大模型应用开发】【项目实战】35.基于A2A协议的智能助手(多智能体)(三)项目实现之天气MCP服务器&票务员MCP服务器&订票MCP服务器

6.天气MCP服务器

mcp_weather_server天气 MCP 服务器,提供 weather_data 表的 SELECT 查询接口,返回 JSON 格式结果

核心功能:

  • 初始化 MySQL 数据库连接

  • 执行 SELECT 查询,返回 JSON 格式结果

  • 格式化日期和数值字段,确保 JSON 序列化兼容

  • 通过 FastAPI 提供 HTTP 接口响应 MCP 工具调用

6.1 格式编码

format.py中包含一个编码器方法和JSON编码器类

目标:定义编码器方法,用于格式化单个对象;自定义 JSON 编码器,处理 MySQL 查询结果中的非标准类型

功能:将 MySQL查询结果中的date、datetime、timedelta 和 Decimal 类型转换为 JSON 兼容的字符串或数值

位置:/utils/format.py

import json
from datetime import date, datetime, timedelta
from decimal import Decimal

def default_encoder(obj): # 定义编码器方法,用于格式化单个对象
if isinstance(obj, datetime): # 检查是否为datetime,返回带时间的格式化字符串
return obj.strftime('%Y-%m-%d %H:%M:%S')
if isinstance(obj, date): # 检查是否为date,返回日期格式化字符串
return obj.strftime('%Y-%m-%d')
if isinstance(obj, timedelta): # 检查是否为timedelta,转换为字符串
return str(obj)
if isinstance(obj, Decimal): # 检查是否为Decimal,转换为浮点数
return float(obj)
return obj # 否则返回原对象

# 定义自定义JSON编码器类,继承自json.JSONEncoder,用于处理非标准类型序列化
class DateEncoder(json.JSONEncoder):
def default(self, obj): # 重写default方法,处理序列化时的默认对象转换
if isinstance(obj, (date, datetime)): # 检查对象是否为date或datetime类型,对于datetime返回带时间的字符串,对于date返回日期字符串
return obj.strftime('%Y-%m-%d %H:%M:%S') if isinstance(obj, datetime) else obj.strftime('%Y-%m-%d')
if isinstance(obj, timedelta): # 检查对象是否为timedelta类型,将时间差转换为字符串
return str(obj)
if isinstance(obj, Decimal): # 检查对象是否为Decimal类型,将Decimal转换为浮点数以兼容JSON
return float(obj)
return super().default(obj) # 对于其他类型,调用父类默认方法

测试:

if __name__ == '__main__':
print(default_encoder(datetime(2025=6, 7, 11, 8, 0)))
print(default_encoder(date(2026, 7, 11)))
print(default_encoder(timedelta(days=1)))
print(default_encoder(Decimal('123.45')))
print('*'*80)

encoder = DateEncoder()
print(encoder.default(datetime(2026, 7, 11, 8, 0)))
print(encoder.default(date(2026, 7, 11)))
print(encoder.default(timedelta(days=1)))
print(encoder.default(Decimal('123.45')))

6.2 WeatherService类

目标:提供天气数据查询服务,响应代理的 SQL 请求

功能:初始化 MySQL 连接,执行 SELECT 查询,格式化结果为 JSON

位置:/mcp_server/mcp_weather_server.py

# -*- coding: utf-8 -*-
"""
这是一个基于 FastMCP 框架的天气查询服务。
它连接 MySQL 数据库,通过接收 SQL 查询语句来返回天气数据。
"""

# 导入必要的库
import mysql.connector # 用于连接和操作 MySQL 数据库
import json # 用于处理 JSON 数据的序列化和反序列化
from datetime import date, datetime, timedelta # 用于处理日期和时间对象
from decimal import Decimal # 用于处理高精度的十进制数,如数据库中的 DECIMAL 类型

# 从自定义模块中导入 FastMCP 类,这是构建 MCP 服务器的核心
from mcp.server.fastmcp import FastMCP

# 从本地配置文件和日志模块中导入配置和日志记录器
from config import Config
from create_logger import logger

# 从自定义工具模块中导入日期编码器和默认编码器,用于格式化特殊数据类型
from utils.format import DateEncoder, default_encoder

# 实例化配置对象,用于获取数据库连接信息(host, user, password, database)
conf = Config()

# 定义天气服务类,用于封装所有与天气数据相关的数据库操作
class WeatherService:
def __init__(self):
# 在类初始化时,使用配置信息建立与 MySQL 数据库的连接
# 这个连接会在整个 WeatherService 实例的生命周期内保持
self.conn = mysql.connector.connect(
host=conf.host,
user=conf.user,
password=conf.password,
database=conf.database
)

# 定义一个执行 SQL 查询的方法
# 参数: sql (str) – 一个完整的 SQL 查询字符串
# 返回: str – 一个 JSON 格式的字符串,包含查询结果或错误信息
def execute_query(self, sql: str) -> str:
try:
# 创建一个游标(cursor),dictionary=True 表示查询结果将以字典列表的形式返回
# 这样可以通过列名(如 row['city'])而不是索引来访问数据
cursor = self.conn.cursor(dictionary=True)

# 执行传入的 SQL 语句
cursor.execute(sql)

# 获取所有查询结果
results = cursor.fetchall()

# 关闭游标,释放资源
cursor.close()

# — 数据格式化 —
# 遍历查询结果中的每一行(每个结果都是一个字典)
for result in results:
# 遍历字典中的每一个键值对
for key, value in result.items():
# 检查值是否是日期、时间或 Decimal 等无法直接被 JSON 序列化的特殊类型
if isinstance(value, (date, datetime, timedelta, Decimal)):
# 如果是,则使用自定义的 default_encoder 函数将其转换为字符串
result[key] = default_encoder(value)

# — 构建并返回 JSON 响应 —
# 根据查询结果是否为空,构建不同的响应结构
# 如果有数据,返回 {"status": "success", "data": […]}
# 如果无数据,返回 {"status": "no_data", "message": "…"}
response_data = {"status": "success", "data": results} if results else {"status": "no_data",
"message": "未找到天气数据,请确认城市和日期。"}

# 将 Python 字典序列化为 JSON 字符串
# cls=DateEncoder: 使用自定义的 JSONEncoder 来处理序列化过程中遇到的其他特殊类型
# ensure_ascii=False: 确保中文字符不被转义,直接显示
return json.dumps(response_data, cls=DateEncoder, ensure_ascii=False)

except Exception as e:
# — 异常处理 —
# 如果在执行过程中发生任何错误(如 SQL 语法错误、数据库连接失败等)
# 使用日志记录器记录错误信息
logger.error(f"天气查询错误: {str(e)}")

# 返回一个包含错误状态的 JSON 字符串
return json.dumps({"status": "error", "message": str(e)}, ensure_ascii=False)

# 用于独立测试 WeatherService 类
# if __name__ == "__main__":
# service = WeatherService()
# sql = "SELECT * FROM weather_data WHERE city='西安' limit 2"
# print(service.execute_query(sql))

# {"status": "success", "data": results}

# {"status": "success",
# "data":
# [
# {"id": 1, "city": "西安", "fx_date": "2026-07-04", "sunrise": "7:21:00",
# "sunset": "17:37:00", "moonrise": "20:21:00", "moonset": "8:41:00",
# "moon_phase": "亏凸月", "moon_phase_icon": "805", "temp_max": 12,
# "temp_min": -2, "icon_day": "100", "text_day": "晴", "icon_night": "150",
# "text_night": "晴", "wind360_day": 225, "wind_dir_day": "西南风", "wind_scale_day": "1-3",
# "wind_speed_day": 3, "wind360_night": 0, "wind_dir_night": "北风", "wind_scale_night": "1-3",
# "wind_speed_night": 16, "precip": 0.0, "uv_index": 3, "humidity": 29, "pressure": 1016,
# "vis": 25, "cloud": 0, "update_time": "2026-07-04 11:13:00"},
# {"id": 2, "city": "西安", "fx_date": "2026-07-05", "sunrise": "7:20:00", "sunset": "17:38:00",
# "moonrise": "21:26:00", "moonset": "9:03:00", "moon_phase": "亏凸月", "moon_phase_icon": "805",
# "temp_max": 3, "temp_min": -7, "icon_day": "101", "text_day": "多云", "icon_night": "151",
# "text_night": "多云", "wind360_day": 0, "wind_dir_day": "北风", "wind_scale_day": "1-3",
# "wind_speed_day": 16, "wind360_night": 45, "wind_dir_night": "东北风", "wind_scale_night": "1-3",
# "wind_speed_night": 3, "precip": 0.0, "uv_index": 2, "humidity": 19, "pressure": 1030, "vis": 25,
# "cloud": 0, "update_time": "2026-07-04 11:13:00"}
# ]
# }

基于 MCP (Model Context Protocol) 协议的天气查询服务。它封装了数据库操作,并通过 HTTP 接口暴露查询能力,允许外部系统(如 AI 模型)通过 SQL 语句查询天气数据

6.3 启动MCP服务器

create_weather_mcp_server()函数

目标:创建并启动天气 MCP 服务器

功能:初始化 FastMCP,注册 query_weather 工具,启动 FastAPI 服务器,监听端口 6001

# 创建并运行 MCP 服务器的主函数
def create_weather_mcp_server():
# 创建一个 FastMCP 服务器实例
weather_mcp = FastMCP(
name="WeatherTools", # 服务器的名称
instructions="天气查询工具,基于 weather_data 表。", # 服务器的描述或指令,供 AI 模型理解其功能
log_level="ERROR", # 设置日志级别,只显示 ERROR 及以上级别的日志
host="127.0.0.1", # 服务器监听的 IP 地址
port=8002 # 服务器监听的端口号
)

# 实例化之前定义的 WeatherService 类,用于处理具体的查询逻辑
service = WeatherService()

# 使用装饰器将一个函数注册为 MCP 服务器的“工具”(tool)
# 这使得 AI 模型可以调用这个函数
@weather_mcp.tool(
name="query_weather", # 工具的名称
# 工具的描述,详细说明了工具的功能和输入参数(一个 SQL 字符串)的格式
description="查询天气数据,输入 SQL,如 'SELECT * FROM weather_data WHERE city = \\"北京\\" AND fx_date = \\"2025-07-30\\"'"
)
def query_weather(sql: str) -> str:
# 当工具被调用时,首先记录一条信息日志
logger.info(f"执行天气查询: {sql}")
# 调用 WeatherService 实例的 execute_query 方法来执行 SQL 并返回结果
return service.execute_query(sql)

# — 服务器启动 —
# 打印服务器信息到控制台
logger.info("=== 天气MCP服务器信息 ===")
logger.info(f"名称: {weather_mcp.name}")
logger.info(f"描述: {weather_mcp.instructions}")

# 运行服务器,使其开始监听请求
try:
print("服务器已启动,请访问 http://127.0.0.1:8002/mcp")
# 使用 streamable-http 作为传输协议启动服务器
weather_mcp.run(transport="streamable-http")
except Exception as e:
# 如果服务器启动失败,打印错误信息
print(f"服务器启动失败: {e}")

# Python 的标准入口点
# 当直接运行此脚本时,下面的代码块会被执行
if __name__ == '__main__':
create_weather_mcp_server()

/mcp_server/mcp_weather_server.py整体代码:

# -*- coding: utf-8 -*-
"""
这是一个基于 FastMCP 框架的天气查询服务。
它连接 MySQL 数据库,通过接收 SQL 查询语句来返回天气数据。
"""

# 导入必要的库
import mysql.connector # 用于连接和操作 MySQL 数据库
import json # 用于处理 JSON 数据的序列化和反序列化
from datetime import date, datetime, timedelta # 用于处理日期和时间对象
from decimal import Decimal # 用于处理高精度的十进制数,如数据库中的 DECIMAL 类型

# 从自定义模块中导入 FastMCP 类,这是构建 MCP 服务器的核心
from mcp.server.fastmcp import FastMCP

# 从本地配置文件和日志模块中导入配置和日志记录器
from config import Config
from create_logger import logger

# 从自定义工具模块中导入日期编码器和默认编码器,用于格式化特殊数据类型
from utils.format import DateEncoder, default_encoder

# 实例化配置对象,用于获取数据库连接信息(host, user, password, database)
conf = Config()

# 定义天气服务类,用于封装所有与天气数据相关的数据库操作
class WeatherService:
def __init__(self):
# 在类初始化时,使用配置信息建立与 MySQL 数据库的连接
# 这个连接会在整个 WeatherService 实例的生命周期内保持
self.conn = mysql.connector.connect(
host=conf.host,
user=conf.user,
password=conf.password,
database=conf.database
)

# 定义一个执行 SQL 查询的方法
# 参数: sql (str) – 一个完整的 SQL 查询字符串
# 返回: str – 一个 JSON 格式的字符串,包含查询结果或错误信息
def execute_query(self, sql: str) -> str:
try:
# 创建一个游标(cursor),dictionary=True 表示查询结果将以字典列表的形式返回
# 这样可以通过列名(如 row['city'])而不是索引来访问数据
cursor = self.conn.cursor(dictionary=True)

# 执行传入的 SQL 语句
cursor.execute(sql)

# 获取所有查询结果
results = cursor.fetchall()

# 关闭游标,释放资源
cursor.close()

# — 数据格式化 —
# 遍历查询结果中的每一行(每个结果都是一个字典)
for result in results:
# 遍历字典中的每一个键值对
for key, value in result.items():
# 检查值是否是日期、时间或 Decimal 等无法直接被 JSON 序列化的特殊类型
if isinstance(value, (date, datetime, timedelta, Decimal)):
# 如果是,则使用自定义的 default_encoder 函数将其转换为字符串
result[key] = default_encoder(value)

# — 构建并返回 JSON 响应 —
# 根据查询结果是否为空,构建不同的响应结构
# 如果有数据,返回 {"status": "success", "data": […]}
# 如果无数据,返回 {"status": "no_data", "message": "…"}
response_data = {"status": "success", "data": results} if results else {"status": "no_data",
"message": "未找到天气数据,请确认城市和日期。"}

# 将 Python 字典序列化为 JSON 字符串
# cls=DateEncoder: 使用自定义的 JSONEncoder 来处理序列化过程中遇到的其他特殊类型
# ensure_ascii=False: 确保中文字符不被转义,直接显示
return json.dumps(response_data, cls=DateEncoder, ensure_ascii=False)

except Exception as e:
# — 异常处理 —
# 如果在执行过程中发生任何错误(如 SQL 语法错误、数据库连接失败等)
# 使用日志记录器记录错误信息
logger.error(f"天气查询错误: {str(e)}")

# 返回一个包含错误状态的 JSON 字符串
return json.dumps({"status": "error", "message": str(e)}, ensure_ascii=False)

# 创建并运行 MCP 服务器的主函数
def create_weather_mcp_server():
# 创建一个 FastMCP 服务器实例
weather_mcp = FastMCP(
name="WeatherTools", # 服务器的名称
instructions="天气查询工具,基于 weather_data 表。", # 服务器的描述或指令,供 AI 模型理解其功能
log_level="ERROR", # 设置日志级别,只显示 ERROR 及以上级别的日志
host="127.0.0.1", # 服务器监听的 IP 地址
port=8002 # 服务器监听的端口号
)

# 实例化之前定义的 WeatherService 类,用于处理具体的查询逻辑
service = WeatherService()

# 使用装饰器将一个函数注册为 MCP 服务器的“工具”(tool)
# 这使得 AI 模型可以调用这个函数
@weather_mcp.tool(
name="query_weather", # 工具的名称
# 工具的描述,详细说明了工具的功能和输入参数(一个 SQL 字符串)的格式
description="查询天气数据,输入 SQL,如 'SELECT * FROM weather_data WHERE city = \\"北京\\" AND fx_date = \\"2025-07-30\\"'"
)
def query_weather(sql: str) -> str:
# 当工具被调用时,首先记录一条信息日志
logger.info(f"执行天气查询: {sql}")
# 调用 WeatherService 实例的 execute_query 方法来执行 SQL 并返回结果
return service.execute_query(sql)

# — 服务器启动 —
# 打印服务器信息到控制台
logger.info("=== 天气MCP服务器信息 ===")
logger.info(f"名称: {weather_mcp.name}")
logger.info(f"描述: {weather_mcp.instructions}")

# 运行服务器,使其开始监听请求
try:
print("服务器已启动,请访问 http://127.0.0.1:8002/mcp")
# 使用 streamable-http 作为传输协议启动服务器
weather_mcp.run(transport="streamable-http")
except Exception as e:
# 如果服务器启动失败,打印错误信息
print(f"服务器启动失败: {e}")

# Python 的标准入口点
# 当直接运行此脚本时,下面的代码块会被执行
if __name__ == '__main__':
create_weather_mcp_server()

# 用于独立测试 WeatherService 类
# if __name__ == "__main__":
# service = WeatherService()
# sql = "SELECT * FROM weather_data WHERE city='西安' limit 2"
# print(service.execute_query(sql))

# {"status": "success", "data": results}

# {"status": "success",
# "data":
# [
# {"id": 1, "city": "西安", "fx_date": "2026-07-04", "sunrise": "7:21:00",
# "sunset": "17:37:00", "moonrise": "20:21:00", "moonset": "8:41:00",
# "moon_phase": "亏凸月", "moon_phase_icon": "805", "temp_max": 12,
# "temp_min": -2, "icon_day": "100", "text_day": "晴", "icon_night": "150",
# "text_night": "晴", "wind360_day": 225, "wind_dir_day": "西南风", "wind_scale_day": "1-3",
# "wind_speed_day": 3, "wind360_night": 0, "wind_dir_night": "北风", "wind_scale_night": "1-3",
# "wind_speed_night": 16, "precip": 0.0, "uv_index": 3, "humidity": 29, "pressure": 1016,
# "vis": 25, "cloud": 0, "update_time": "2026-07-04 11:13:00"},
# {"id": 2, "city": "西安", "fx_date": "2026-07-05", "sunrise": "7:20:00", "sunset": "17:38:00",
# "moonrise": "21:26:00", "moonset": "9:03:00", "moon_phase": "亏凸月", "moon_phase_icon": "805",
# "temp_max": 3, "temp_min": -7, "icon_day": "101", "text_day": "多云", "icon_night": "151",
# "text_night": "多云", "wind360_day": 0, "wind_dir_day": "北风", "wind_scale_day": "1-3",
# "wind_speed_day": 16, "wind360_night": 45, "wind_dir_night": "东北风", "wind_scale_night": "1-3",
# "wind_speed_night": 3, "precip": 0.0, "uv_index": 2, "humidity": 19, "pressure": 1030, "vis": 25,
# "cloud": 0, "update_time": "2026-07-04 11:13:00"}
# ]
# }

客户端测试:

位置:/test/test_weather_mcp_server.py

异步 Python 脚本,用于连接并测试之前创建的 WeatherTools MCP 服务器,。它模拟了一个 AI 客户端,通过 MCP 协议与服务器通信,并调用其提供的 query_weather 工具

# -*- coding: utf-8 -*-
"""
这是一个用于测试 WeatherTools MCP 服务器的客户端脚本。
它使用异步 I/O 来连接服务器、获取工具列表并执行查询。
"""

# 导入必要的库
import asyncio # 用于编写异步代码,处理并发操作
import json # 用于解析从服务器返回的 JSON 字符串

# 从 langchain_mcp_adapters 库中导入工具,用于从 MCP 会话中加载可用工具
from langchain_mcp_adapters.tools import load_mcp_tools

# 从 mcp 库中导入核心客户端类
from mcp import ClientSession
# 从 mcp 库中导入用于建立 streamable-http 连接的客户端函数
from mcp.client.streamable_http import streamablehttp_client

# 定义要连接的 MCP 服务器地址
# 这与您之前启动的服务器地址和端口一致
server_url = "http://127.0.0.1:8002/mcp"

# 定义一个异步函数来执行测试
async def test_weather_mcp():
try:
# — 1. 建立与服务器的连接 —
# 使用 streamablehttp_client 作为异步上下文管理器来连接到服务器
# 它会返回一对读写通道 (read, write) 和一个 get_session_id 函数(此处用 _ 忽略)
async with streamablehttp_client(server_url) as (read, write, _):

# — 2. 初始化 MCP 会话 —
# 使用上一步获取的读写通道创建一个 ClientSession 实例
async with ClientSession(read, write) as session:
try:
# 向服务器发送初始化请求,建立 MCP 协议会话
await session.initialize()
print("会话初始化成功,可以开始调用工具。")

# — 3. 获取服务器提供的工具列表 —
# 调用 load_mcp_tools 函数,它会自动向服务器查询并加载所有可用的工具
tools = await load_mcp_tools(session)
# 打印工具列表,用于验证连接和工具发现是否成功
print(f"tools–>{tools}")

# — 4. 测试调用工具:查询指定日期天气 —
# 构造一个 SQL 查询字符串,作为工具的输入参数
sql = "SELECT * FROM weather_data WHERE city = '西安' AND fx_date = '2026-07-04'"

# 使用 session.call_tool 方法调用服务器上名为 "query_weather" 的工具
# 第一个参数是工具名,第二个参数是一个字典,包含工具所需的参数
result = await session.call_tool("query_weather", {"sql": sql})

# 打印原始返回结果,用于调试
print(11111, result)
# 注意:result 不是一个简单的字符串,而是一个包含 TextContent 和 structuredContent 的对象

# 检查 result 是否为字符串类型
print(22222, isinstance(result, str))
# 输出为 False,证实 result 是一个复杂对象

# 尝试将 result 解析为 JSON
# 这里的逻辑有误,因为 result 本身不是 JSON 字符串,所以 json.loads(result) 会失败
# 正确的做法是提取 result.content[0].text 或 result.structuredContent['result']
result_data = json.loads(result) if isinstance(result, str) else result

# 打印最终处理后的结果数据
print(f"指定日期天气结果:{result_data}")
# 此时 result_data 仍然是原始的复杂对象,包含了服务器返回的 JSON 字符串

# — 5. (已注释) 测试调用工具:查询未来3天天气 —
# 以下是另一个测试用例,用于查询一个日期范围内的天气
# sql_range = "SELECT * FROM weather_data WHERE city = '西安' AND fx_date BETWEEN '2026-07-04' AND '2026-07-06'"
# result_range = await session.call_tool("query_weather", {"sql": sql_range})
# result_range_data = json.loads(result_range) if isinstance(result_range, str) else result_range
# print(f"天气范围查询结果:{result_range_data}")

except Exception as e:
# 捕获并打印在会话内部发生的任何错误(如工具调用失败)
print(f"天气 MCP 测试出错:{str(e)}")

except Exception as e:
# 捕获并打印在连接或会话初始化阶段发生的错误
print(f"连接或会话初始化时发生错误: {e}")
print("请确认服务端脚本已启动并运行在 http://127.0.0.1:8002/mcp")

# Python 的标准入口点
# 使用 asyncio.run() 来运行顶层的异步函数
if __name__ == "__main__":
asyncio.run(test_weather_mcp())

7.票务MCP服务器

mcp_ticket_server.py:票务 MCP 服务器,提供 train_tickets、flight_tickets 和 concert_tickets 表的 SELECT 查询接口,返回 JSON 格式结果

核心功能:

  • 初始化 MySQL 数据库连接
  • 执行 SELECT 查询,返回 JSON 格式结果
  • 格式化日期和数值字段,确保 JSON 序列化兼容
  • 通过 FastAPI 提供 HTTP 接口,响应 MCP 工具调用

7.1 TicketService类

目标:提供票务数据查询服务,响应代理的 SQL 请求

功能:初始化 MySQL 连接,执行 SELECT 查询,格式化结果为 JSON

位置:/query_data/query_ticket.py

票务查询服务(TicketService),它的核心逻辑与之前的天气服务非常相似,主要是负责连接数据库、执行 SQL 并返回 JSON 格式的结果

# 导入必要的库
import mysql.connector # 用于连接和操作 MySQL 数据库
import json # 用于处理 JSON 数据的序列化和反序列化
from datetime import date, datetime, timedelta # 用于处理日期和时间对象
from decimal import Decimal # 用于处理高精度的十进制数(如数据库中的金额、DECIMAL 类型)

# 从 mcp 库中导入 FastMCP 类,用于构建 MCP 服务器(虽然在此代码片段中未使用,但已导入)
from mcp.server.fastmcp import FastMCP

# 从本地自定义模块中导入配置、日志和格式化工具
from config import Config
from create_logger import logger
from utils.format import DateEncoder, default_encoder

# 实例化配置对象,用于获取数据库连接信息(host, user, password, database)
conf = Config()

# 定义票务服务类,封装所有与票务数据相关的数据库操作
class TicketService:
def __init__(self):
# 在类初始化时,使用配置信息建立与 MySQL 数据库的连接
# 这个连接会在整个 TicketService 实例的生命周期内保持
self.conn = mysql.connector.connect(
host=conf.host,
user=conf.user,
password=conf.password,
database=conf.database
)

# 定义一个执行 SQL 查询的方法
# 参数: sql (str) – 一个完整的 SQL 查询字符串
# 返回: str – 一个 JSON 格式的字符串,包含查询结果或错误信息
def execute_query(self, sql: str) -> str:
try:
# 创建一个游标(cursor),dictionary=True 表示查询结果将以字典列表的形式返回
# 这样可以通过列名(如 row['city'])而不是索引来访问数据
cursor = self.conn.cursor(dictionary=True)

# 执行传入的 SQL 语句
cursor.execute(sql)

# 获取所有查询结果
results = cursor.fetchall()

# 关闭游标,释放资源
cursor.close()

# — 数据格式化 —
# 遍历查询结果中的每一行(每个结果都是一个字典)
for result in results:
# 遍历字典中的每一个键值对
for key, value in result.items():
# 检查值是否是日期、时间或 Decimal 等无法直接被 JSON 序列化的特殊类型
if isinstance(value, (date, datetime, timedelta, Decimal)):
# 如果是,则使用自定义的 default_encoder 函数将其转换为字符串
result[key] = default_encoder(value)

# — 构建并返回 JSON 响应 —
# 根据查询结果是否为空,构建不同的响应结构
# 如果有数据,返回 {"status": "success", "data": […]}
# 如果无数据,返回 {"status": "no_data", "message": "未找到票务数据,请确认查询条件。"}
response_data = {"status": "success", "data": results} if results else {
"status": "no_data",
"message": "未找到票务数据,请确认查询条件。"
}

# 将 Python 字典序列化为 JSON 字符串
# cls=DateEncoder: 使用自定义的 JSONEncoder 来处理序列化过程中遇到的其他特殊类型
# ensure_ascii=False: 确保中文字符不被转义,直接显示
return json.dumps(response_data, cls=DateEncoder, ensure_ascii=False)

except Exception as e:
# — 异常处理 —
# 如果在执行过程中发生任何错误(如 SQL 语法错误、数据库连接失败等)
# 使用日志记录器记录错误信息
logger.error(f"票务查询错误: {str(e)}")

# 返回一个包含错误状态的 JSON 字符串
return json.dumps({"status": "error", "message": str(e)}, ensure_ascii=False)

总结:代码是一个标准的数据库查询封装层。它接收 SQL 语句,处理了查询过程中可能出现的特殊数据类型(如日期、金额),统一了成功、无数据和异常三种情况的返回格式,最终输出标准的 JSON 字符串

7.2 启动MCP 服务器

create_ticket_mcp_server()函数

目标:创建并启动票务 MCP 服务器

功能:初始化 FastMCP,注册 query_tickets 工具,启动 FastAPI 服务器,监听端口 6002

位置: /mcp_server/mcp_ticket_server.py

# -*- coding: utf-8 -*-
"""
票务 MCP 服务器启动脚本。
此脚本负责创建一个 MCP 服务器,并将 TicketService 的查询功能注册为工具。
"""

# 导入必要的库
import mysql.connector
import json
from datetime import date, datetime, timedelta
from decimal import Decimal

# 从 mcp 库中导入 FastMCP,这是构建 MCP 服务器的核心类
from mcp.server.fastmcp import FastMCP

# 导入自定义模块
from config import Config
from create_logger import logger
from utils.format import DateEncoder, default_encoder
# 导入之前定义的票务服务类
from query_data.query_ticket import TicketService

# 实例化配置对象,用于加载数据库连接信息
conf = Config()

# 定义一个函数来创建并运行票务 MCP 服务器
def create_ticket_mcp_server():
# — 1. 创建 FastMCP 服务器实例 —
# name: 服务器的名称,客户端会发现并使用这个名字
# instructions: 服务器的指令或描述,告诉 AI 这个工具是做什么的
# log_level: 设置日志级别为 ERROR,只记录错误信息
# host & port: 服务器监听的地址和端口
ticket_mcp = FastMCP(
name="TicketTools",
instructions="票务查询工具,基于 train_tickets, flight_tickets, concert_tickets 表。只支持查询。",
log_level="ERROR",
host="127.0.0.1",
port=8001
)

# — 2. 实例化业务逻辑服务对象 —
# 创建一个 TicketService 实例,用于处理具体的数据库查询逻辑
service = TicketService()

# — 3. 注册 MCP 工具 —
# 使用 @ticket_mcp.tool 装饰器将一个普通函数注册为 MCP 工具
@ticket_mcp.tool(
name="query_tickets", # 工具的名称,客户端将通过此名称调用
description="查询票务数据,输入 SQL,如 'SELECT * FROM train_tickets WHERE departure_city = \\"北京\\" AND arrival_city = \\"上海\\"'" # 工具的描述,帮助 AI 理解如何使用它
)
def query_tickets(sql: str) -> str:
# 这是一个普通的 Python 函数,但被装饰器包装后成为了一个 MCP 工具
# 当客户端调用 'query_tickets' 工具时,这个函数就会被执行

# 记录日志,方便调试和追踪
logger.info(f"执行票务查询: {sql}")

# 调用 TicketService 实例的 execute_query 方法来执行 SQL 并获取结果
# 然后将结果(一个 JSON 字符串)返回给客户端
return service.execute_query(sql)

# — 4. 打印服务器信息 —
logger.info("=== 票务MCP服务器信息 ===")
logger.info(f"名称: {ticket_mcp.name}")
logger.info(f"描述: {ticket_mcp.instructions}")

# — 5. 启动服务器 —
try:
print("服务器已启动,请访问 http://127.0.0.1:8001/mcp")
# 运行服务器,并指定传输方式为 "streamable-http"
# 这是一种基于 HTTP 长连接的通信方式,适合流式数据传输
ticket_mcp.run(transport="streamable-http")
except Exception as e:
# 捕获并打印启动过程中可能出现的任何错误
print(f"服务器启动失败: {e}")

# 调用函数,启动服务器
create_ticket_mcp_server()

MCP 服务器启动脚本,它整合了之前定义的 TicketService 类,将票务查询功能通过 MCP 协议暴露出去

总结: 代码是整个票务查询服务的入口点。它完成了以下几件事:

  • 配置服务器:使用 FastMCP 类创建了一个名为 TicketTools 的服务器实例,并设置了其名称、描述、监听地址和端口
  • 集成业务逻辑:实例化了 TicketService 类,将数据库操作的能力引入
  • 暴露工具:通过 @ticket_mcp.tool 装饰器,将 query_tickets 函数注册为一个可供外部调用的 MCP 工具,。这个工具接收一个 SQL 字符串作为参数,并返回查询结果的 JSON 字符串
  • 启动服务:最后,调用 ticket_mcp.run() 方法,以 streamable-http 模式启动服务器,使其开始监听客户端的连接和请求
  • 客户端测试:

    位置:/test/test_ticket_mcp_server.py

    import asyncio
    import json

    # 从 langchain_mcp_adapters 库中导入工具,用于自动发现和加载 MCP 服务器提供的工具
    from langchain_mcp_adapters.tools import load_mcp_tools

    # 从 mcp 库中导入核心客户端类,用于建立会话
    from mcp import ClientSession

    # 从 mcp 库中导入用于建立 streamable-http 连接的客户端函数
    from mcp.client.streamable_http import streamablehttp_client

    # 定义要连接的 MCP 服务器地址,与您提供的票务服务器地址一致
    server_url = "http://127.0.0.1:8001/mcp"

    async def test_ticket_mcp():
    try:
    # — 1. 建立与服务器的连接 —
    # 使用 streamablehttp_client 作为异步上下文管理器来连接到服务器
    # 它会返回一对读写通道 (read, write) 和一个 get_session_id 函数(此处用 _ 忽略)
    async with streamablehttp_client(server_url) as (read, write, _):

    # — 2. 初始化 MCP 会话 —
    # 使用上一步获取的读写通道创建一个 ClientSession 实例
    async with ClientSession(read, write) as session:
    try:
    # 向服务器发送初始化请求,建立 MCP 协议会话
    await session.initialize()
    print("会话初始化成功,可以开始调用工具。")

    # — 3. 获取服务器提供的工具列表 —
    # 调用 load_mcp_tools 函数,它会自动向服务器查询并加载所有可用的工具
    tools = await load_mcp_tools(session)
    print(f"tools–>{tools}")

    # — 4. 测试调用工具:查询机票 —
    # 构造一个 SQL 查询字符串,用于查询从上海到北京的公务舱机票
    sql_flights = "SELECT * FROM flight_tickets WHERE departure_city = '上海' AND arrival_city = '北京' AND DATE(departure_time) = '2025-10-28' AND cabin_type = '公务舱'"
    # 使用 session.call_tool 方法调用服务器上名为 "query_tickets" 的工具
    result_flights = await session.call_tool("query_tickets", {"sql": sql_flights})
    # 处理返回结果,尝试将其解析为 JSON
    result_flights_data = json.loads(result_flights) if isinstance(result_flights, str) else result_flights
    print(f"机票查询结果:{result_flights_data}")

    # — 5. 测试调用工具:查询火车票 —
    # 构造一个 SQL 查询字符串,用于查询从北京到上海的二等座火车票
    sql_trains = "SELECT * FROM train_tickets WHERE departure_city = '北京' AND arrival_city = '上海' AND DATE(departure_time) = '2025-10-22' AND seat_type = '二等座'"
    # 调用工具并处理结果
    result_trains = await session.call_tool("query_tickets", {"sql": sql_trains})
    result_trains_data = json.loads(result_trains) if isinstance(result_trains, str) else result_trains
    print(f"火车票查询结果:{result_trains_data}")

    # — 6. 测试调用工具:查询演唱会票 —
    # 构造一个 SQL 查询字符串,用于查询北京刀郎演唱会的看台票
    sql_concerts = "SELECT * FROM concert_tickets WHERE city = '北京' AND artist = '刀郎' AND DATE(start_time) = '2025-10-31' AND ticket_type = '看台'"
    # 调用工具并处理结果
    result_concerts = await session.call_tool("query_tickets", {"sql": sql_concerts})
    result_concerts_data = json.loads(result_concerts) if isinstance(result_concerts, str) else result_concerts
    print(f"演唱会票查询结果:{result_concerts_data}")

    except Exception as e:
    # 捕获并打印在会话内部发生的任何错误(如工具调用失败)
    print(f"票务 MCP 测试出错:{str(e)}")

    except Exception as e:
    # 捕获并打印在连接或会话初始化阶段发生的错误
    print(f"连接或会话初始化时发生错误: {e}")
    print("请确认服务端脚本已启动并运行在 http://127.0.0.1:8001/mcp")

    # Python 的标准入口点
    # 使用 asyncio.run() 来运行顶层的异步函数
    if __name__ == "__main__":
    asyncio.run(test_ticket_mcp())

    测试 票务 MCP 服务器 的客户端脚本: 它通过异步方式连接到运行在 8001 端口的服务器,并依次测试了机票、火车票和演唱会门票的查询功能

    总结: 代码是一个完整的 MCP 客户端测试脚本,它模拟了 AI 代理与票务服务交互的全过程:

  • 建立连接:使用 streamablehttp_client 连接到运行在 8001 端口的票务 MCP 服务器
  • 初始化会话:通过 ClientSession 与服务器建立一个 MCP 协议会话
  • 发现工具:使用 load_mcp_tools 自动获取服务器上所有可用的工具列表
  • 调用工具:脚本依次构造了三个不同的 SQL 查询,分别用于测试机票、火车票和演唱会门票的查询功能,并通过 session.call_tool 方法调用远程的 query_tickets 工具
  • 处理结果:脚本接收服务器返回的结果,并尝试将其解析为 JSON 格式后打印出来
  • 8.订票MCP服务器

    mcp_order_server.py:订票 MCP 服务器,通过调用API完成火车票,飞机票和演唱会票的预定

    核心功能:

    • 火车票预定、飞机票预定、演出票预定
    • 通过 FastAPI 提供 HTTP 接口,响应 MCP 工具调用

    位置:/mcp_server/mcp_order_server.py

    from mcp.server.fastmcp import FastMCP

    from config import Config
    from create_logger import logger

    # 加载配置并创建日志记录器
    conf = Config()

    # — 1. 创建 FastMCP 实例 —
    # 初始化一个名为 "OrderTools" 的 MCP 服务器实例
    # – name: 服务器的名称
    # – instructions: 服务器的功能描述,会提供给客户端
    # – log_level: 日志级别,设置为 "ERROR"
    # – host: 服务器监听的主机地址
    # – port: 服务器监听的端口号,此处为 8003
    order_mcp = FastMCP(
    name="OrderTools",
    instructions="票务预定工具,通过调用API完成火车票、飞机票和演唱会票的预定。",
    log_level="ERROR",
    host="127.0.0.1",
    port=8003
    )

    # — 2. 定义工具:预定火车票 —
    # 使用 @order_mcp.tool 装饰器将一个普通函数注册为 MCP 工具
    @order_mcp.tool(
    name="order_train", # 工具的唯一标识名
    description="根据时间、车次、座位类型、数量预定火车票" # 工具的功能描述
    )
    def order_train(departure_date: str, train_number: str, seat_type: str, number: int) -> str:
    '''
    Args:
    departure_date (str): 出发日期,如 '2025-10-30'
    train_number (str): 火车车次,如 'G346'
    seat_type (str): 座位类型,如 '二等座'
    number (int): 订购张数
    '''
    # 记录订票日志
    logger.info(f"正在订购火车票: {departure_date}, {train_number}, {seat_type}, {number}")
    logger.info(f"恭喜,火车票预定成功!")
    # 返回预定成功的消息
    return "恭喜,火车票预定成功!"

    # — 3. 定义工具:预定飞机票 —
    @order_mcp.tool(
    name="order_flight",
    description="根据时间、班次、座位类型、数量预定飞机票"
    )
    def order_flight(departure_date: str, flight_number: str, seat_type: str, number: int) -> str:
    '''
    Args:
    departure_date (str): 出发日期,如 '2025-10-30'
    flight_number (str): 飞机班次,如 'CA6557'
    seat_type (str): 座位类型,如 '经济舱'
    number (int): 订购张数
    '''
    # 记录订票日志
    logger.info(f"正在订购飞机票: {departure_date}, {flight_number}, {seat_type}, {number}")
    logger.info(f"恭喜,飞机票预定成功!")
    # 返回预定成功的消息
    return "恭喜,飞机票预定成功!"

    # — 4. 定义工具:预定演唱会票 —
    @order_mcp.tool(
    name="order_concert",
    description="根据时间、明星、场地、座位类型、数量预定演出票"
    )
    def order_concert(start_date: str, aritist: str, venue: str, seat_type: str, number: int) -> str:
    '''
    Args:
    start_date (str): 开始日期,如 '2025-10-30'
    aritist (str): 明星,如 '刀郎'
    venue (str): 场地,如 '上海体育馆'
    seat_type (str): 座位类型,如 '看台'
    number (int): 订购张数
    '''
    # 记录订票日志
    logger.info(f"正在订购演出票: {start_date}, {aritist}, {venue}, {seat_type}, {number}")
    logger.info(f"恭喜,演出票预定成功!")
    # 返回预定成功的消息
    return "恭喜,演出票预定成功!"

    # — 5. 创建并运行服务器 —
    def create_order_mcp_server():
    # 打印服务器信息到日志
    logger.info("=== 票务预定MCP服务器信息 ===")
    logger.info(f"名称: {order_mcp.name}")
    logger.info(f"描述: {order_mcp.instructions}")

    # 运行服务器
    try:
    print("服务器已启动,请访问 http://127.0.0.1:8003/mcp")
    # 启动服务器,并指定使用 "streamable-http" 作为传输协议
    order_mcp.run(transport="streamable-http")
    except Exception as e:
    print(f"服务器启动失败: {e}")

    # Python 的标准入口点
    if __name__ == "__main__":
    # 调用函数启动服务器
    create_order_mcp_server()

    基于 FastMCP 框架构建的 票务预定 MCP 服务器。它定义了三个核心工具,分别用于预定火车票、飞机票和演唱会门票,并通过 streamable-http 协议在 8003 端口对外提供服务

    客户端测试:

    位置:/test/test_order_mcp_server.py

    import asyncio
    import json

    # 从 LangChain 导入构建 Agent 所需的核心组件
    from langchain.agents import create_tool_calling_agent, AgentExecutor
    from langchain_core.prompts import ChatPromptTemplate
    from langchain_mcp_adapters.tools import load_mcp_tools
    from langchain_openai import ChatOpenAI

    # 从 MCP 库导入建立客户端会话所需的组件
    from mcp import ClientSession
    from mcp.client.streamable_http import streamablehttp_client

    # 从项目配置中导入配置和日志记录器
    from config import Config
    from create_logger import logger

    conf = Config()

    # — 1. 初始化大语言模型 (LLM) —
    # 使用 ChatOpenAI 类初始化一个 LLM 实例
    # 模型的名称、API 地址和密钥都从配置文件中加载
    llm = ChatOpenAI(
    model=conf.model_name,
    base_url=conf.base_url,
    api_key=conf.api_key,
    temperature=0.1 # 设置较低的温度,使模型输出更确定
    )

    # — 2. 定义异步的票务预定主函数 —
    async def order_tickets(query):
    try:
    # — 2.1 建立与 MCP 服务器的连接 —
    # 使用 streamablehttp_client 连接到票务预定服务器
    async with streamablehttp_client("http://127.0.0.1:8003/mcp") as (read, write, _):
    # 使用读写通道创建一个 MCP 客户端会话
    async with ClientSession(read, write) as session:
    try:
    # 初始化会话
    await session.initialize()

    # — 2.2 动态加载工具 —
    # 从已连接的 MCP 服务器自动发现并加载所有可用的工具
    # 这些工具(如 order_train, order_flight)将作为 Agent 的“手”
    tools = await load_mcp_tools(session)

    # — 2.3 创建 Agent 的提示词模板 —
    # 定义 Agent 的行为准则和对话结构
    prompt = ChatPromptTemplate.from_messages([
    ("system",
    "你是一个票务预定助手,能够调用工具来完成火车票、飞机票或演出票的预定。你需要仔细分析工具需要的参数,然后从用户提供的信息中提取信息。如果用户提供的信息不足以提取到调用工具所有必要参数,则向用户追问,以获取该信息。不能自己编撰参数。"),
    ("human", "{input}"),
    ("placeholder", "{agent_scratchpad}"),
    ])

    # — 2.4 构建并执行 Agent —
    # 1. 创建工具调用 Agent:将 LLM、工具和提示词模板组合在一起
    agent = create_tool_calling_agent(llm, tools, prompt)

    # 2. 创建 Agent 执行器:负责运行 Agent 的逻辑,verbose=True 会打印详细过程
    agent_executor = AgentExecutor(agent=agent, tools=tools, verbose=True)

    # 3. 调用 Agent:将用户的查询传入,并等待执行结果
    response = await agent_executor.ainvoke({"input": query})

    # 返回 Agent 的最终输出
    return response['output']
    except Exception as e:
    # 捕获并记录会话内部的错误
    logger.info(f"票务 MCP 测试出错:{str(e)}")
    return f"票务 MCP 查询出错:{str(e)}"
    except Exception as e:
    # 捕获并记录连接或会话初始化的错误
    logger.error(f"连接或会话初始化时发生错误: {e}")
    return "连接或会话初始化时发生错误"

    # — 3. 主程序入口 —
    if __name__ == "__main__":
    # 创建一个简单的命令行交互循环
    while True:
    query = input("请输入查询:")
    if query == "exit":
    break
    # 运行异步的 order_tickets 函数并打印结果
    print(asyncio.run(order_tickets(query)))

    基于 LangChain 框架构建的 AI 票务预定助手。它作为客户端,连接到运行在 8003 端口的票务预定 MCP 服务器,通过大语言模型(LLM)理解用户意图,并自动调用相应的工具来完成火车票、飞机票或演唱会门票的预定

    【上一篇】【AI大模型应用开发】【项目实战】34.基于A2A协议的智能助手(多智能体)(二)项目实现之项目架构图&配置模块&数据模块

    【下一篇】【AI大模型应用开发】【项目实战】36.基于A2A协议的智能助手(多智能体)(四)项目实现之天气Agent服务器&票务员Agent服务器&订票Agent服务器

    赞(0)
    未经允许不得转载:网硕互联帮助中心 » 【AI大模型应用开发】【项目实战】35.基于A2A协议的智能助手(多智能体)(三)项目实现之天气MCP服务器&票务员MCP服务器&订票MCP服务器
    分享到: 更多 (0)

    评论 抢沙发

    评论前必须登录!