Token导航 LogoToken导航TokenDH.com
SSB MCP Server logo
数据服务stdio官方级别未说明来源级核验

SSB MCP Server

MCP Server

提供对SQL流构建器(SSB)的全面访问和管理功能,支持多种认证方式和80+管理工具,适用于数据工程和实时流处理场景。

工具数

80

提示词数

0

GitHub Stars

4

资源数

0
实时分析PythonClaudeClaude DesktopClaude

安装说明

本站只整理中文说明和来源信息,不托管安装包,也不代用户安装。

作者 / 组织

BrooksIan

提供方

BrooksIan

最后核验

2026/5/17 20:29

运行时

Python

快速接入

先看主来源和安装命令,再打开仓库或文档;下面只保留这个条目的关键接入事实。

命令预览

python3 -m venv .venv

详细介绍

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)

  1. 启动SSB服务:
   git clone https://github.com/your-org/ssb-mcp-server.git
   cd ssb-mcp-server
   docker-compose up -d
  1. 配置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"
       }
     }
   }
  1. 重新启动克劳德桌面 并开始与您的SSB流进行交互!

用于CDP SSB部署

您的SSB API基本URL通常为:

https:///ssb/api/v1

从CDP UI获取Knox JWT令牌,并将其与以下配置一起使用。

Knox Token Generation

设置

选项1:克劳德桌面(本地)

  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 .
  1. 配置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"
          }
        }
      }
    }
  1. 重新启动克劳德桌面 并开始询问有关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_USERSSB用户名(例如。, admin)
SSB_PASSWORDSSB密码(例如。, admin)
SSB_READONLY只读模式(默认: false)
TIMEOUT_SECONDSHTTP超时(秒)(默认值: 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_BUNDLECA证书包的路径
SSB_READONLY只读模式(默认: true)
TIMEOUT_SECONDSHTTP超时(秒)(默认值: 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流生成器的全面访问。以下是一些功能的可视化示例:

标准品牌盒子版本:

可用工具

🔧 高级作业管理

  • 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服务器可以从现有的物化视图中检索数据,但不能以编程方式创建新的物化图。要创建物化视图,请执行以下操作:

  1. 在SSB UI中导航到您的作业
  2. 转到“物化视图”部分
  3. 通过UI配置和创建MV
  4. 使用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) -流开始端点返回404
  • stop_stream(stream_name) -流停止端点返回404
  • execute_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_usernamemve_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表需要两个步骤:

  1. 创建数据源:使用 create_kafka_table() 创建数据源配置
  2. 在目录中注册:使用 register_kafka_table() 使表可查询

自动配准:The register_kafka_table() 功能:

  • 使用DDL在Flink目录中注册表(默认值: ssb.ssb_default,回落到 default_catalog.ssb_default)
  • 可配置的目录和数据库参数,用于灵活的命名空间控制
  • 当请求的目录不可用时,自动回退目录
  • 默认情况下,在与现有表(如NVDA)相同的数据库中创建表
  • 自动添加 ssb_ 表名前缀(可配置)
  • 根据主题数据自动创建架构
  • 通过检查正确的数据库上下文验证表是否可用于查询
  • 返回成功注册的确认信息,包括完整的表名、目录和数据库信息

命名规范:

  • 默认:桌子得到 ssb_ 自动前缀(例如。, user_eventsssb_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部署中以只读模式运行,在独立部署中启用写。要更改此设置,请执行以下操作:

  1. SSB_READONLY=false (允许写入)或 SSB_READONLY=true (只读)
  2. 重新启动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

故障排除

常见问题

  1. “未经授权”的错误:检查您的身份验证凭据

- 对于直接SSB:验证 SSB_USERSSB_PASSWORD - 对于诺克斯:验证 KNOX_TOKENKNOX_USER/KNOX_PASSWORD

  1. “连接被拒绝”错误:确保SSB服务正在运行

- 检查 docker-compose ps 关于服务状态 - 验证docker-compose.yml中的端口映射

  1. “没有可用的样本数据”:工作可能需要时间来生成数据

- 检查作业状态 get_job_status(job_id) - 验证作业是否正在运行并具有示例配置 - 尝试使用 sample_all_messages=True 用于全面采样

  1. 作业重启失败:如果作业重新启动失败

- 使用 restart_job_with_sampling() 而不是 execute_job() - 检查作业是否处于允许重新启动的状态 - 如果无法重新启动,则创建新作业 - 注: start_stream()stop_stream() 方法不起作用(404错误)

  1. SSL证书错误:用于诺克斯部署

- 集 KNOX_VERIFY_SSL=false 用于自签名证书 - 或提供适当的CA捆绑包 KNOX_CA_BUNDLE

  1. Kafka表创建错误:如果表创建失败

- 验证是否仅使用本地kafka连接器(强制) - 检查Kafka主题是否存在并且可以访问 - 确保引导服务器配置正确 - 使用 validate_kafka_connector() 检查连接器有效性

  1. 虚拟表不可查询:创建Kafka表后

- 重要:创建数据源不会自动使其可用于查询 - 数据源需要通过SSB UI在Flink目录中手动注册 - 使用 SHOW TABLES; 查看哪些表实际上可用于查询 - 仅显示在中的表 SHOW TABLES; 可以通过SQL查询

  1. 配置未更新:更新后 config/cloud_ssb_config.json

- 重要:环境变量优先于JSON配置 - 对于JSON配置:更新 config/cloud_ssb_config.json 并重新启动MCP服务器 - 对于环境变量: export KNOX_TOKEN="new_token" 并重新启动 - 检查您的部署正在使用哪种方法

  1. 物化视图访问错误:访问MVE API时

- MVE API需要基本身份验证(用户名:密码) - SSB API使用承载令牌身份验证 - 使用不同的凭据访问MVE API - 在访问其MV之前,请检查作业是否正在运行

  1. 空表列表:跑步时 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

目录标签

目录标签

实时分析PythonClaude流数据处理本地部署SQL流构建器大数据管理ApacheKnox集成

支持客户端

Claude DesktopClaude

接入字段

传输方式(transport,传输协议)

stdio

鉴权方式(authType,认证方式)

token

运行时(runtime,运行环境)

Python

工具数量(toolCount,工具数)

80

资源数量(resourceCount,资源数)

0

提示词数量(promptCount,提示词数)

0

权限和风险

stdiotoken部署方式未说明

接入前请确认传输方式、认证方式和部署位置,并根据实际工具能力限制访问范围。

安装前确认

不要直接授予不必要的文件、网络或账号权限;先核对安装命令和配置内容。

来源信息

继续浏览同类 MCP