SSB MCP服务器
SSB Home Dashboard *主SSB界面显示了带有可用流和作业的主仪表板,其特点是SSB MCP服务器与突出品牌的集成。*
模型上下文协议服务器提供对SQL流生成器(SSB)的全面访问,支持直接SSB访问和Apache Knox集成。
适用于独立SSB部署和Cloudera数据平台(CDP)SSB部署 -通过Claude Desktop提供对SSB功能的完全访问。
特性
- 多种身份验证方法:
- 直接SSB身份验证:独立SSB部署的基本身份验证 - 诺克斯集成:CDP部署的承载令牌、Cookie和密码令牌
- 默认情况下为只读 -SSB流和配置的安全探索
- 全面的SSB API覆盖范围 随着 80+MCP工具 对于完整的SSB管理:
- 高级作业管理:事件历史记录、状态管理、作业复制、数据源克隆 - 监测和诊断:系统运行状况、性能计数器、SQL分析 - 增强的表管理:详细表格、层次结构、验证、创建 - 连接器和格式管理:数据格式、连接器详细信息、JAR信息 - 用户和项目管理:设置、项目、用户信息、项目创建 - API密钥管理:密钥生命周期管理、创建、删除、详细信息 - 环境管理:环境切换、配置、创建 - 同步和配置:项目导出/导入、同步管理、验证 - UDF管理:UDF生命周期、执行、工件、自定义函数 - 流管理:列出、创建、更新、删除、启动、停止SQL流 - 查询执行:执行SQL查询并通过采样获得实时结果 - 示例数据访问:从正在运行的作业中检索流数据样本 - 作业管理:监视作业状态、获取作业详细信息并管理作业生命周期 - 架构发现:探索表模式和可用表 - 功能管理:列出并检查用户定义的函数 - 连接器管理:探索可用的连接器 - Kafka集成:列出并检查Kafka主题 - 集群监控:获取群集信息和运行状况 - 性能指标:监控流性能和指标
快速开始
用于独立SSB部署(Docker Compose)
- 启动SSB服务:
git clone https://github.com/your-org/ssb-mcp-server.git
cd ssb-mcp-server
docker-compose up -d- 配置Claude桌面 -编辑
~/Library/Application Support/Claude/claude_desktop_config.json:
{
"mcpServers": {
"ssb-mcp-server": {
"command": "/FULL/PATH/TO/SSB-MCP-Server/run_mcp_server.sh",
"args": [],
"cwd": "/FULL/PATH/TO/SSB-MCP-Server"
}
}
}- 重新启动克劳德桌面 并开始与您的SSB流进行交互!
用于CDP SSB部署
您的SSB API基本URL通常为:
https:///ssb/api/v1从CDP UI获取Knox JWT令牌,并将其与以下配置一起使用。
设置
选项1:克劳德桌面(本地)
- 克隆并安装:
git clone https://github.com/your-org/ssb-mcp-server.git
cd ssb-mcp-server
python3 -m venv .venv
source .venv/bin/activate
pip install -e .- 配置Claude桌面 -编辑
~/Library/Application Support/Claude/claude_desktop_config.json:
{
"mcpServers": {
"ssb-mcp-server": {
"command": "/FULL/PATH/TO/SSB-MCP-Server/.venv/bin/python",
"args": [
"-m",
"ssb_mcp_server.server"
],
"env": {
"MCP_TRANSPORT": "stdio",
"SSB_API_BASE": "https://ssb-gateway.yourshere.cloudera.site/ssb/api/v1",
"KNOX_TOKEN": "",
"SSB_READONLY": "true"
}
}
}
}- 重新启动克劳德桌面 并开始询问有关SSB流的问题!
选项2:Cloudera AI代理工作室
要与Cloudera AI Agent Studio一起使用,请使用环境变量进行配置:
{
"mcpServers": {
"ssb-mcp-server": {
"command": "uvx",
"args": [
"--from",
"git+https://github.com/your-org/ssb-mcp-server@main",
"run-server"
],
"env": {
"MCP_TRANSPORT": "stdio",
"SSB_BASE_URL": "https://your-ssb-host:443",
"SSB_TENANT": "your-tenant-name",
"KNOX_TOKEN": "your_jwt_token_here",
"MVE_USERNAME": "your_mve_username",
"MVE_PASSWORD": "your_mve_password",
"SSB_READONLY": "false"
}
}
}
}配置模板:使用 config/cloudera_ai_agent_studio_template.json 作为一个起点。
Cloudera AI Agent Studio的环境变量:
export SSB_BASE_URL="https://your-ssb-host:443"
export SSB_TENANT="your-tenant-name"
export KNOX_TOKEN="your_jwt_token_here"
export MVE_USERNAME="your_mve_username"
export MVE_PASSWORD="your_mve_password"
export SSB_READONLY="false"配置选项
配置可以通过环境变量或JSON配置文件完成。环境变量优先于JSON配置。
JSON配置文件
MCP服务器可以从以下位置加载配置 config/cloud_ssb_config.json:
{
"cloud_ssb": {
"base_url": "https://your-ssb-host:443",
"tenant": "your-tenant-name",
"jwt_token": "your_jwt_token_here",
"ssb_readonly": false,
"knox_verify_ssl": true,
"http_timeout_seconds": 60,
"http_max_retries": 3,
"http_rate_limit_rps": 5,
"mve_username": "your_mve_username",
"mve_password": "your_mve_password"
}
}环境变量
所有配置也可以通过环境变量完成:
直接SSB身份验证(独立)
| 变量 | 必填 | 描述 |
|---|---|---|
SSB_API_BASE | 是 | 完整的SSB API URL(例如。, http://localhost:18121) |
SSB_USER | 是 | SSB用户名(例如。, admin) |
SSB_PASSWORD | 是 | SSB密码(例如。, admin) |
SSB_READONLY | 否 | 只读模式(默认: false) |
TIMEOUT_SECONDS | 否 | HTTP超时(秒)(默认值: 30) |
Knox身份验证(CDP)
| 变量 | 必填 | 描述 |
|---|---|---|
SSB_BASE_URL | 是\* | SSB主机的基本URL(例如。, https://your-host:443) |
SSB_TENANT | 是\* | 租户名称(例如。, your-tenant-name) |
KNOX_TOKEN | 是\* | 用于身份验证的Knox JWT令牌 |
KNOX_COOKIE | 否 | 替代方案:提供完整的cookie字符串而不是令牌 |
KNOX_PASSCODE_TOKEN | 否 | 替代方案:Knox密码令牌(自动兑换为JWT) |
KNOX_USER | 否 | 基本身份验证的Knox用户名 |
KNOX_PASSWORD | 否 | 基本身份验证的Knox密码 |
KNOX_VERIFY_SSL | 否 | 验证SSL证书(默认值: true) |
KNOX_CA_BUNDLE | 否 | CA证书包的路径 |
SSB_READONLY | 否 | 只读模式(默认: true) |
TIMEOUT_SECONDS | 否 | HTTP超时(秒)(默认值: 30) |
MVE_USERNAME | 没有用于基本身份验证的 | MVE API用户名 |
MVE_PASSWORD | 没有用于基本身份验证的 | MVE API密码 |
传统配置(仍受支持)
| 变量 | 必填 | 描述 |
|---|---|---|
KNOX_GATEWAY_URL | 否 | 遗留:完整的诺克斯网关URL |
SSB_API_BASE | 否 | 旧版:完整的SSB API URL |
MVE_API_BASE | 否 | 旧版:完整的MVE API URL |
\*要么 SSB_API_BASE (直接)或 KNOX_GATEWAY_URL (适用于诺克斯)为必填项
示例功能
SSB MCP服务器通过Claude Desktop提供对SQL流生成器的全面访问。以下是一些功能的可视化示例:
标准品牌盒子版本:
- SSB Home Branded -顶部标有黄色文字
可用工具
🔧 高级作业管理
get_job_events(job_id)-获取详细的作业事件历史记录和时间线get_job_state(job_id)-获取全面的作业状态信息get_job_mv_endpoints(job_id)-获取作业的物化视图端点create_job_mv_endpoint(job_id, mv_config)-创建或更新物化视图端点copy_job(job_id)-复制现有作业copy_data_source(data_source_id)-克隆数据源
📊 监测和诊断
get_diagnostic_counters()-获取系统性能计数器和诊断get_heartbeat()-检查系统运行状况和连接性analyze_sql(sql_query)-无需执行即可分析SQL查询(语法、性能)
🗂️ 增强的表管理
list_tables_detailed()-获取全面的表格信息get_table_tree()-获取按目录组织的分层表结构validate_data_source(data_source_config)-验证数据源配置create_table_detailed(table_config)-创建具有完整配置的表get_table_details(table_id)-获取特定表格的详细信息
🔌 连接器和格式管理
list_data_formats()-列出所有可用的数据格式get_data_format_details(format_id)-获取特定数据格式的详细信息create_data_format(format_config)-创建新的数据格式get_connector_jar(connector_type)-获取连接器JAR信息get_connector_type_details(connector_type)-获取详细的连接器类型信息get_connector_details(connector_id)-获取详细的连接器信息
👤 用户和项目管理
get_user_settings()-获取用户偏好和设置update_user_settings(settings)-更新用户配置list_projects()-列出可用项目get_project_details(project_id)-获取项目信息create_project(project_config)-创建新项目get_user_info()-获取当前用户信息
🔑 API密钥管理
list_api_keys()-列出用户API密钥create_api_key(key_config)-创建新的API密钥delete_api_key(key_id)-删除API密钥get_api_key_details(key_id)-获取API关键信息
🌍 环境管理
list_environments()-列出可用环境activate_environment(env_id)-激活/切换到环境get_environment_details(env_id)-获取环境配置create_environment(env_config)-创建新环境deactivate_environment()-停用当前环境
🔄 同步和配置
get_sync_config()-获取同步配置update_sync_config(config)-更新同步配置delete_sync_config()-删除同步配置validate_sync_config(project)-验证项目的同步配置export_project(project)-导出项目配置import_project(project, config)-导入项目配置
📈 UDF管理
list_udfs_detailed()-获取全面的UDF信息run_udf(udf_id, parameters)-执行UDF函数get_udf_artifacts()-获取UDF工件和依赖关系create_udf(udf_config)-创建自定义自定义自定义项update_udf(udf_id, udf_config)-更新UDF配置get_udf_details(udf_id)-获取详细的UDF信息get_udf_artifact_details(artifact_id)-获取UDF工件详细信息get_udf_artifact_by_type(artifact_type)-按类型获取UDF工件
流管理
list_streams()-列出所有SQL流(作业)get_stream(stream_name)-获取特定流的详细信息create_stream(stream_name, sql_query, description?)-创建新流(写入模式)update_stream(stream_name, sql_query, description?)-更新现有流delete_stream(stream_name)-删除流start_stream(stream_name)-启动流stop_stream(stream_name)-停下一条小溪
查询执行和示例数据
execute_query(sql_query, limit?)-执行SQL查询并创建SSB作业execute_query_with_sampling(sql_query, sample_interval, sample_count, window_size, sample_all_messages)-使用自定义采样执行查询get_job_status(job_id)-获取特定SSB作业的状态get_job_sample(sample_id)-从作业执行中获取示例数据get_job_sample_by_id(job_id)-按作业ID从作业中获取示例数据list_jobs_with_samples()-列出所有作业及其示例信息
作业管理与控制
stop_job(job_id, savepoint)-停止特定的SSB作业execute_job(job_id, sql_query)-使用新的SQL执行/重新启动作业restart_job_with_sampling(job_id, sql_query, sample_interval, sample_all_messages)-使用采样选项重新启动作业configure_sampling(sample_id, sample_interval, sample_count, window_size, sample_all_messages)-配置采样参数
数据源和模式
list_tables()-列出所有可用表(数据源)get_table_schema(table_name)-获取表架构信息create_kafka_table(table_name, topic, kafka_connector_type, bootstrap_servers, format_type, scan_startup_mode, additional_properties?)-使用本地kafka强制创建新表register_kafka_table(table_name, topic, schema_fields?, use_ssb_prefix?, catalog?, database?)-在Flink目录中注册Kafka表(使其可查询)validate_kafka_connector(kafka_connector_type)-验证连接器类型是否为本地kafka
功能和连接器
list_udfs()-列出用户定义的函数get_udf(udf_name)-获取UDF详细信息list_connectors()-列出可用连接器get_connector(connector_name)-获取连接器详细信息
Kafka集成
list_topics()-列出Kafka主题get_topic(topic_name)-获取主题详细信息
群集管理
get_cluster_info()-获取集群信息get_cluster_health()-获取群集运行状况get_ssb_info()-获取SSB版本和系统信息
示例用法
配置后,您可以向Claude提出以下问题:
基本信息
- “我运行的是哪个版本的SSB?”
- “列出我的所有SQL流”
- “哪些表可用于查询?”
- “有哪些可用的连接器?”
- “列出所有Kafka主题”
- “集群健康状况如何?”
Claude Show Tables *Claude可以显示SSB环境中的所有可用表,包括内置表和自定义表。*
查询执行和数据访问
- “执行此查询:从NVDA中选择\*”
- “使用示例所有消息创建作业:SELECT\*FROM NVDA”
- “显示作业1234的状态”
Claude List Jobs *Claude可以列出所有正在运行的SSB作业,显示其状态、创建时间和详细信息。*
- “从作业1234获取示例数据”
- “列出所有作业及其示例信息”
- “显示ReadNVDA作业的最新数据”
作业管理与控制
- “停止作业1234”
- “使用所有消息样本重新启动作业1234”
- “使用自定义采样(间隔500ms)重新启动作业1234”
- “配置作业1234的采样以采样所有消息”
Claude Created Job *Claude可以通过执行SQL查询创建新的SSB作业,并具有完整的作业管理功能。*
流管理
- “使用以下SQL创建一个名为‘sales_analysis’的新流:SELECT\*FROM sales amount>1000”
- “显示'user_events'流的详细信息”
- “我的‘sales_stream’的状态如何?”
Kafka表管理
- “从主题‘用户事件’创建一个名为‘user_events’的本地Kafka表”
- “在Flink目录中注册一个Kafka表,使其可查询”
- “创建JSON格式的本地Kafka表”
Claude Created Table *Claude可以通过适当的模式和连接器配置创建连接到Kafka主题的新虚拟表。*
- “验证'local kafka'是否是有效的连接器类型”
- “为实时数据流创建虚拟表”
高级作业管理
- “显示作业1234的事件历史记录”
- “获取作业1234的详细状态”
- “复制作业1234以创建新作业”
- “克隆表'user_events'的数据源”
- “获取作业1234的物化视图端点”
- “为作业1234创建物化视图端点”
物化视图
- “获取作业1234的物化视图端点”
- “为作业1234创建物化视图端点”
⚠️ 重要限制:必须通过SSB UI界面创建物化视图(MV)。MCP服务器可以从现有的物化视图中检索数据,但不能以编程方式创建新的物化图。要创建物化视图,请执行以下操作:
- 在SSB UI中导航到您的作业
- 转到“物化视图”部分
- 通过UI配置和创建MV
- 使用MCP服务器查询创建的MV数据
⚠️ 已知限制
作业重启功能
有限的重启功能:SSB API提供有限的作业重新启动功能:
✅ 什么有效:
stop_job(job_id, savepoint=True)-使用保存点停止作业restart_job_with_sampling(job_id, sql_query, ...)-使用更新的SQL创建新作业PUT /jobs/{id}-更新作业配置
❌ 什么不起作用:
- 通过专用重启端点直接重启作业
start_stream(stream_name)-流开始端点返回404stop_stream(stream_name)-流停止端点返回404execute_job(job_id, sql_query)-数据库连接问题
推荐的重启策略:
# For restarting a job with new SQL
def restart_siem_job(job_id, new_sql):
try:
# Stop the existing job
client.stop_job(job_id, savepoint=True)
# Create new job with updated SQL
result = client.restart_job_with_sampling(job_id, new_sql)
return result
except Exception as e:
print(f"Restart failed: {e}")虚拟表和系统目录
有限表发现:SSB环境对传统表的支持有限:
✅ 可用内容:
- 系统目录:2个目录(
default_catalog,ssb) - 函数:206个内置数据处理函数
- 物化视图:可通过MVE API访问(非SQL)
❌ 什么不可用:
- 传统数据库表(SHOW tables返回空)
- 标准信息_模式查询
- 对物化视图的直接SQL访问
- 用于元数据发现的系统表
数据访问方法:
- 直接sql:仅限于功能和基本查询
- MVE API:对于物化视图数据(需要基本身份验证)
- 作业执行:用于创建和管理数据流
身份验证和配置
配置加载:MCP服务器现在支持JSON配置文件和环境变量:
✅ JSON配置:更新 config/cloud_ssb_config.json 并重新启动MCP服务器 ✅ 环境变量:设置环境变量(优先于JSON) ✅ MVE API证书:配置 mve_username 和 mve_password JSON或环境变量
⚠️ 重要:对于Cloudera AI Agent Studio,使用环境变量:
export KNOX_GATEWAY_URL="your_gateway_url"
export KNOX_TOKEN="your_jwt_token"
export SSB_API_BASE="your_api_base"
export MVE_USERNAME="your_mve_username"
export MVE_PASSWORD="your_mve_password"令牌到期:JWT令牌过期,需要手动刷新。MCP服务器不会自动刷新令牌。
云环境限制
端点可用性:某些端点在云环境中可能不可用:
❌ 云限制:
list_projects()-可能返回“不支持请求方法'GET'”get_heartbeat()-可能返回空响应create_stream()和execute_query()-可能超时500个错误SHOW JOBS-云环境不支持
✅ 解决方法:
- 使用
list_streams()而不是list_projects() - 在云环境中跳过心跳检查
- 使用重试逻辑优雅地处理超时
- 使用作业管理工具而不是SHOW命令
物化视图引擎(MVE)API
单独身份验证:MVE API需要与主SSB API不同的身份验证:
身份验证要求:
- API SSB:承载令牌身份验证
- MVE API:基本身份验证(用户名:密码)
- 凭证:必须单独提供MVE访问权限
访问模式:
# MVE API access requires Basic Auth
credentials = 'username:password'
encoded_credentials = base64.b64encode(credentials.encode()).decode()
headers = {'Authorization': f'Basic {encoded_credentials}'}监测和诊断
- “检查系统心跳和健康状况”
- “显示诊断计数器”
- “分析此SQL查询的性能:SELECT\*FROM NVDA WHERE close>100”
- “当前系统性能如何?”
增强的表管理
- “显示所有表的详细信息”
- “按目录获取分层表结构”
- “验证此数据源配置”
- “创建具有完整配置的新表”
- “获取有关表'user_events'的详细信息”
Claude Get Table Info *Claude可以提供有关特定表的详细信息,包括它们的模式和配置。*
用户和项目管理
- “显示我的用户设置和首选项”
- “更新我的用户设置以启用黑暗模式”
- “列出所有可用项目”
- “创建一个名为‘分析’的新项目”
- “获取项目'ffffffff'的详细信息”
- “显示我的用户信息”
API密钥管理
- “列出我的所有API密钥”
- “为外部访问创建新的API密钥”
- “删除API密钥'key123'”
- “获取有关API密钥'key123'的详细信息”
环境管理
- “列出所有可用环境”
- “切换到环境‘生产’”
- 创建一个名为“staging”的新环境
- “获取有关环境‘dev’的详细信息”
- “停用当前环境”
同步和配置
- “显示当前同步配置”
- “更新Git集成的同步配置”
- “导出项目‘分析’配置”
- “从Git导入项目配置”
- 验证项目“测试”的同步配置
UDF管理
- “列出所有用户定义的函数及其详细信息”
- “使用参数运行UDF'custom_aggregate'”
- “为数据转换创建新的UDF”
- “更新UDF‘my_function’配置”
- “获取UDF工件和依赖项”
示例数据示例
MCP服务器可以使用不同的采样模式检索实时流数据样本:
Claude Job Sample *Claude可以从正在运行的作业中检索实时示例数据,显示实际的流数据。*
定期采样(默认):
{
"records": [
{
"___open": "185.0919",
"___high": "185.1200",
"___low": "184.9400",
"___close": "184.9700",
"___volume": "61884",
"eventTimestamp": "2025-10-08T18:34:10.915Z"
}
],
"job_status": "RUNNING",
"end_of_samples": false,
"message": "Retrieved 1 sample records"
}示例所有消息模式:
{
"sampling_mode": "sample_all_messages",
"sample_interval": 0,
"sample_count": 10000,
"window_size": 10000,
"message": "Job created with comprehensive sampling enabled"
}SQL查询功能
- 自动半导体处理:所有SQL查询都会自动以分号终止
- 灵活采样:在定期采样或采样所有消息之间进行选择
- 作业控制:使用不同配置启动、停止和重新启动作业
- 实时数据:从正在运行的作业中访问流数据样本
高级功能
示例所有消息
要进行全面的数据采样,请使用 sample_all_messages=True 选项:
# Create job with sample all messages
execute_query_with_sampling("SELECT * FROM NVDA", sample_all_messages=True)
# Restart job with sample all messages
restart_job_with_sampling(1234, "SELECT * FROM NVDA", sample_all_messages=True)配置:
sample_interval: 0(立即取样)sample_count: 10000(捕获所有消息的计数很高)window_size: 10000(用于全面采样的大窗口)
自定义采样配置
根据您的特定需求微调采样行为:
# Custom sampling with 500ms interval
execute_query_with_sampling("SELECT * FROM NVDA",
sample_interval=500,
sample_count=500,
window_size=500)
# Configure existing job sampling
configure_sampling("sample_id",
sample_interval=200,
sample_count=1000,
window_size=1000)作业管理
完整的作业生命周期管理:
# Stop job with savepoint
stop_job(1234, savepoint=True)
# Restart job with new SQL
execute_job(1234, "SELECT * FROM NEW_TABLE")
# Restart with sampling options
restart_job_with_sampling(1234, "SELECT * FROM NVDA",
sample_interval=1000,
sample_all_messages=False)Kafka表创建
创建仅限于本地kafka连接器的表:
# Step 1: Create data source (creates configuration)
create_kafka_table("user_events", "user-events") # Uses local-kafka by default
# Step 2: Register table in Flink catalog (makes it queryable)
register_kafka_table("user_events", "user-events") # Creates ssb_user_events in ssb.ssb_default (falls back to default_catalog.ssb_default)
# Advanced: Custom schema registration
custom_schema = [
{"name": "id", "type": "STRING"},
{"name": "name", "type": "STRING"},
{"name": "timestamp", "type": "TIMESTAMP"}
]
register_kafka_table("custom_table", "custom-topic", custom_schema) # Creates ssb_custom_table
# Without ssb_ prefix
register_kafka_table("raw_data", "raw-topic", use_ssb_prefix=False) # Creates raw_data (no prefix)
# Custom catalog and database
register_kafka_table("custom_table", "custom-topic", catalog="default_catalog", database="default_database")
# Local Kafka with custom settings
create_kafka_table("local_data", "local-topic", "local-kafka", "localhost:9092", "json", "earliest-offset")
register_kafka_table("local_data", "local-topic") # Creates ssb_local_data
# Validate connector types
validate_kafka_connector("local-kafka") # Returns validation details
validate_kafka_connector("kafka") # Returns error - only local-kafka allowed两步流程:创建Kafka表需要两个步骤:
- 创建数据源:使用
create_kafka_table()创建数据源配置 - 在目录中注册:使用
register_kafka_table()使表可查询
自动配准:The register_kafka_table() 功能:
- 使用DDL在Flink目录中注册表(默认值:
ssb.ssb_default,回落到default_catalog.ssb_default) - 可配置的目录和数据库参数,用于灵活的命名空间控制
- 当请求的目录不可用时,自动回退目录
- 默认情况下,在与现有表(如NVDA)相同的数据库中创建表
- 自动添加
ssb_表名前缀(可配置) - 根据主题数据自动创建架构
- 通过检查正确的数据库上下文验证表是否可用于查询
- 返回成功注册的确认信息,包括完整的表名、目录和数据库信息
命名规范:
- 默认:桌子得到
ssb_自动前缀(例如。,user_events→ssb_user_events) - 以(权力)否决:使用
use_ssb_prefix=False禁用前缀 - 现有的:表已以开头
ssb_未修改
命名空间配置:
- 默认:表创建于
ssb.ssb_default命名空间(回退到default_catalog.ssb_default如果ssb目录不可用) - 自定义目录:使用
catalog用于指定不同目录的参数 - 自定义数据库:使用
database用于指定不同数据库的参数 - 自动回退:系统自动回退到
default_catalog如果请求的目录不可用 - 完全控制:目录和数据库都是可配置的,以实现最大的灵活性
验证:
- 使用
SHOW TABLES;确认表可供查询 - 表创建于
default_catalog.ssb_default命名空间(与NVDA相同) - 使用完整命名空间(
default_catalog.ssb_default.TABLE_NAME)或使用以下命令切换数据库上下文USE default_catalog.ssb_default; - 所有虚拟Kafka表都与现有表位于同一位置,便于查询
支持的连接器:
local-kafka-本地Kafka连接器(虚拟表的唯一选项)
支持的格式:
json-JSON格式(默认)csv-CSV格式avro-Apache Avro格式- 自定义格式字符串
Docker编写设置
该存储库包括一个完整的Docker Compose设置,用于本地开发和测试:
服务包括
- PostgreSQL:SSB元数据数据库
- 卡夫卡:消息流平台
- 弗林克:流处理引擎
- NiFi:数据流管理
- Qdrant:矢量数据库
- SSB SSE:SQL流生成器流式SQL引擎
- SSB MVE:SQL流生成器物化视图引擎
- 阿帕奇诺克斯:安全访问网关(可选)
启动环境
# Start all services
docker-compose up -d
# Check service status
docker-compose ps
# View logs
docker-compose logs -f ssb-sse接入点
- SSB SSE: http://localhost:18121
- SSB MVE: http://localhost:18131
- Flink作业经理: http://localhost:8081
- NiFi: http://localhost:8080
- 诺克斯门户: https://localhost:8444(如果启用)
写入操作
默认情况下,服务器在CDP部署中以只读模式运行,在独立部署中启用写。要更改此设置,请执行以下操作:
- 集
SSB_READONLY=false(允许写入)或SSB_READONLY=true(只读) - 重新启动MCP服务器
写入操作包括:
- 创建、更新和删除流
- 执行创建作业的SQL查询
- 管理作业生命周期(启动、停止、重新启动)
- 配置采样参数
- 作业控制和管理
- 仅创建Kafka表(强制验证)
综合能力
SSB MCP服务器现在提供 80+MCP工具 覆盖 SSB API的80%+,使其成为Claude Desktop提供的最全面的SSB管理平台。
📊 覆盖率统计
- MCP工具总数:80+(高于33)
- API覆盖范围:80%+(高于20%)
- 功能类别:15(高于6)
- 可用端点:67岁以上(高于15岁)
🎯 核心能力
完成SSB管理
- 作业生命周期:创建、监视、控制、复制和管理作业
- 数据管理:表、模式、验证和层次结构组织
- 系统监控:健康检查、诊断和性能跟踪
- 用户管理:设置、项目、环境和API密钥
- DevOps集成:同步、导出/导入和配置管理
高级功能
- 实时采样:灵活的数据采样,带有“采样所有消息”选项
- SQL分析:不执行查询分析以优化性能
- 物化视图:创建和管理物化视图端点
- 自定义UDF:用户定义的功能管理和执行
- 环境控制:具有切换功能的多环境支持
- 项目管理:具有导出/导入功能的完整项目生命周期
企业级就绪
- 安全:API密钥管理和用户身份验证
- 监控:全面的系统健康和性能跟踪
- 可扩展性:支持多个项目和环境
- 集成:Git同步、配置管理和DevOps工作流
- 灵活性:可配置的目录、数据库和命名约定
🚀 用例
数据工程师
- 流处理作业管理和监控
- 实时数据采样和分析
- 表模式管理和验证
- 性能优化和故障排除
DevOps工程师
- 环境管理和配置
- 项目导出/导入和版本控制
- 系统监控和健康检查
- API密钥管理和安全
数据科学家
- 自定义UDF开发和执行
- 数据格式管理和验证
- 查询分析与优化
- 实时数据探索
平台管理员
- 用户和项目管理
- 系统诊断和监控
- 连接器和格式管理
- 同步配置和验证
测试
SSB MCP服务器包括一个全面的测试套件。看 测试/README.md 详细的测试文档,包括:
- 快速功能测试
- 涵盖所有80+MCP工具的全面测试套件
- 云SSB测试协议
- 测试配置和最佳实践
- 详细的测试结果和分析
快速入门:
cd Testing && python run_tests.py --quick故障排除
常见问题
- “未经授权”的错误:检查您的身份验证凭据
- 对于直接SSB:验证 SSB_USER 和 SSB_PASSWORD - 对于诺克斯:验证 KNOX_TOKEN 或 KNOX_USER/KNOX_PASSWORD
- “连接被拒绝”错误:确保SSB服务正在运行
- 检查 docker-compose ps 关于服务状态 - 验证docker-compose.yml中的端口映射
- “没有可用的样本数据”:工作可能需要时间来生成数据
- 检查作业状态 get_job_status(job_id) - 验证作业是否正在运行并具有示例配置 - 尝试使用 sample_all_messages=True 用于全面采样
- 作业重启失败:如果作业重新启动失败
- 使用 restart_job_with_sampling() 而不是 execute_job() - 检查作业是否处于允许重新启动的状态 - 如果无法重新启动,则创建新作业 - 注: start_stream() 和 stop_stream() 方法不起作用(404错误)
- SSL证书错误:用于诺克斯部署
- 集 KNOX_VERIFY_SSL=false 用于自签名证书 - 或提供适当的CA捆绑包 KNOX_CA_BUNDLE
- Kafka表创建错误:如果表创建失败
- 验证是否仅使用本地kafka连接器(强制) - 检查Kafka主题是否存在并且可以访问 - 确保引导服务器配置正确 - 使用 validate_kafka_connector() 检查连接器有效性
- 虚拟表不可查询:创建Kafka表后
- 重要:创建数据源不会自动使其可用于查询 - 数据源需要通过SSB UI在Flink目录中手动注册 - 使用 SHOW TABLES; 查看哪些表实际上可用于查询 - 仅显示在中的表 SHOW TABLES; 可以通过SQL查询
- 配置未更新:更新后
config/cloud_ssb_config.json
- 重要:环境变量优先于JSON配置 - 对于JSON配置:更新 config/cloud_ssb_config.json 并重新启动MCP服务器 - 对于环境变量: export KNOX_TOKEN="new_token" 并重新启动 - 检查您的部署正在使用哪种方法
- 物化视图访问错误:访问MVE API时
- MVE API需要基本身份验证(用户名:密码) - SSB API使用承载令牌身份验证 - 使用不同的凭据访问MVE API - 在访问其MV之前,请检查作业是否正在运行
- 空表列表:跑步时
SHOW TABLES
- SSB环境对传统表的支持有限 - 使用 list_streams() 查看可用的工作 - 通过MVE API而不是SQL访问物化视图 - 专注于基于作业的数据处理,而不是表查询
调试模式
通过设置环境变量启用调试日志记录:
export MCP_LOG_LEVEL=DEBUG安全
- 所有敏感数据(密码、令牌、机密)都会在响应中自动编辑
- 大型收藏被截断,以防止淹没LLM
- 默认情况下,CDP部署启用只读模式,以防止意外修改
- 直接SSB身份验证使用HTTP上的基本身份验证(适用于本地开发)
- SQL查询会自动通过适当的分号终止进行清理
- Kafka表创建强制使用仅限本地Kafka的连接器以确保数据安全
摘要
SSB MCP服务器现在是 综合管理平台 对于SQL流生成器,通过以下方式为Claude Desktop提供对几乎所有SSB功能的访问 80+MCP工具.
🎯 您将获得:
- 完成SSB控制:管理作业、表、用户、项目和环境
- 高级监控:系统运行状况、诊断和性能跟踪
- 实时数据:灵活的采样和流数据访问
- 企业功能:API密钥、同步、导出/导入和多环境支持
- 开发者工具:UDF管理、SQL分析和连接器详细信息
- DevOps集成:项目管理、配置同步和Git工作流
🚀 主要优势:
- 80%+API覆盖率:几乎可以访问所有SSB功能
- 80+MCP工具:为每个用例提供全面的工具集
- 15个功能类别:有组织、可发现的能力
- 企业级就绪:安全性、监控和可扩展性功能
- 用户友好:通过Claude Desktop进行自然语言交互
- 灵活的:支持独立部署和CDP部署
📈 非常适合:
- 数据工程师:流处理、作业管理、实时分析
- DevOps团队:环境管理、监控、配置同步
- 数据科学家:自定义UDF、查询分析、数据探索
- 平台管理员:用户管理、系统监控、安全
SSB MCP服务器将Claude Desktop转换为强大的SSB管理界面,使您能够与整个SQL流生成器环境进行自然语言交互! 🎉
许可证
Apache许可证2.0
