ErlVectorDB
在Erlang/OTP中实现的高性能MCP(模型上下文协议)向量数据库,旨在充分利用Actor模型和OTP监督树进行容错、并发向量操作。
特性
- OTP原生架构:基于gen_server和监督树构建,可实现最大的容错能力
- 分布式聚类:具有自动复制和故障转移功能的水平扩展
- OAuth 2.1身份验证:使用客户端凭据和刷新令牌流进行安全身份验证
- REST API:完整的REST API和MCP协议,用于web集成
- 矢量压缩:多种压缩算法可降低存储需求
- MCP协议支持:用于AI/ML集成的完整模型上下文协议实现
- 永久存储:使用ETS实现基于DETS的持久化,以实现快速内存操作
- 备份和恢复:具有JSON导出/导入的完整备份/还原功能
- 并行矢量运算:利用Erlang的轻量级进程进行并行操作
- 多种距离度量:余弦相似度、欧几里德距离、曼哈顿距离
- 动态店铺管理:在运行时创建和管理多个矢量存储
- 热代码重新加载:在不停机的情况下更新矢量操作
建筑
erlvectordb_sup (Main Supervisor)
├── cluster_manager (Distributed Clustering)
├── oauth_server (OAuth 2.1 Authentication Server)
├── oauth_http_handler (OAuth HTTP Endpoints)
├── rest_api_server (REST API Server)
├── vector_store_sup (Dynamic Supervisor for Vector Stores)
│ ├── vector_store (Gen Server per store)
│ │ └── vector_persistence (DETS + ETS + Compression)
│ └── vector_store (Gen Server per store)
│ └── vector_persistence (DETS + ETS + Compression)
├── mcp_server (MCP Protocol Handler with OAuth)
└── vector_index_manager (Index Management)快速开始
🚀 ErlVectorDB新手? 看看 快速入门指南 只需5分钟的设置!
先决条件
- Erlang/OTP 24+
- 钢筋3
建筑
rebar3 compile跑步
# Quick start with automatic setup
./start-local.sh
# Or manually
rebar3 shell管理服务器
# Check server status
./check-status.sh
# Stop the server and cleanup
./stop-server.sh
# Run automated test suite
./test_server.sh测试
# Run automated test suite
./test_server.sh
# Run Common Test suites
rebar3 ct基本用法
% Start the application
erlvectordb:start().
% Register OAuth client (optional - default admin client is created)
erlvectordb:register_oauth_client(>, >, #{
scopes => [>, >]
}).
% Get OAuth access token
{ok, TokenResponse} = erlvectordb:get_oauth_token(>, >, [>, >]).
AccessToken = maps:get(access_token, TokenResponse).
% Create a regular vector store
erlvectordb:create_store(my_store).
% Create a distributed vector store (if clustering enabled)
erlvectordb:create_distributed_store(distributed_store, #{replication_factor => 2}).
% Insert vectors (with automatic compression if enabled)
Vector1 = [1.0, 2.0, 3.0],
erlvectordb:insert(my_store, >, Vector1, #{title => "Document 1"}).
% Insert with explicit compression
erlvectordb:insert_compressed(my_store, >, [2.0, 3.0, 4.0], #{title => "Document 2"}).
% Search for similar vectors
QueryVector = [1.1, 2.1, 3.1],
{ok, Results} = erlvectordb:search(my_store, QueryVector, 5).
% Get store statistics
{ok, Stats} = erlvectordb:get_stats(my_store).
% Benchmark compression algorithms
Algorithms = [quantization_8bit, quantization_4bit, zlib_compression],
BenchmarkResults = erlvectordb:benchmark_compression([1.0, 2.0, 3.0, 4.0], Algorithms).
% Clustering operations
erlvectordb:join_cluster('other_node@hostname').
{ok, ClusterStatus} = erlvectordb:get_cluster_status().分布式聚类
设置群集
在配置中启用群集:
{cluster_enabled, true},
{node_name, 'erlvectordb@node1.example.com'},
{cluster_cookie, my_secure_cookie},
{replication_factor, 2}集群操作
% Join an existing cluster
erlvectordb:join_cluster('erlvectordb@node2.example.com').
% Create distributed stores
erlvectordb:create_distributed_store(my_distributed_store, #{
replication_factor => 3
}).
% Check cluster status
{ok, Status} = erlvectordb:get_cluster_status().
% Leave cluster
erlvectordb:leave_cluster().自动功能
- 复制:矢量在节点间自动复制
- 故障转移:节点停机时自动故障转移
- 负载平衡:跨可用副本分布的查询
- 一致性:最终符合冲突解决
矢量压缩
支持的算法
- 8位量化:将精度降低到每维8位
- 4位量化:将精度降低到每维4位
- PCA压缩:主成分分析降维
- Zlib压缩:通用压缩
- 产品量化:大矢量的子矢量量化
压缩使用情况
% Enable compression globally
application:set_env(erlvectordb, compression_enabled, true).
application:set_env(erlvectordb, compression_algorithm, quantization_8bit).
% Compress specific vectors
{ok, Compressed} = erlvectordb:compress_vector([1.0, 2.0, 3.0], quantization_8bit).
% Benchmark compression algorithms
Vector = lists:seq(1.0, 100.0, 1.0),
Algorithms = [quantization_8bit, quantization_4bit, zlib_compression],
Results = erlvectordb:benchmark_compression(Vector, Algorithms).
% Results include compression_ratio, compression_time, decompression_time, accuracy_loss压缩优势
- 存储减少:根据算法节省50-90%的存储空间
- 网络效率:更快的复制和备份操作
- 内存使用:减少大型数据集的RAM需求
- 可配置的权衡:压缩比和精度之间的平衡
REST API
该系统提供了一个完整的REST API和MCP协议:
端点
店铺管理:
POST /api/v1/stores-创建商店GET /api/v1/stores-列出商店DELETE /api/v1/stores/{name}-删除商店
矢量操作:
POST /api/v1/stores/{name}/vectors-插入向量POST /api/v1/stores/{name}/search-搜索向量GET /api/v1/stores/{name}/stats-获取店铺统计信息
集群管理:
GET /api/v1/cluster/status-集群状态POST /api/v1/cluster/join-加入集群
REST API示例
# Get OAuth token
curl -X POST http://localhost:8081/oauth/token \
-d "grant_type=client_credentials&client_id=admin&client_secret=admin_secret_2024&scope=read write"
# Create store
curl -X POST http://localhost:8082/api/v1/stores \
-H "Authorization: Bearer YOUR_TOKEN" \
-H "Content-Type: application/json" \
-d '{"name": "my_store"}'
# Insert vector
curl -X POST http://localhost:8082/api/v1/stores/my_store/vectors \
-H "Authorization: Bearer YOUR_TOKEN" \
-H "Content-Type: application/json" \
-d '{"id": "doc1", "vector": [1.0, 2.0, 3.0], "metadata": {"title": "Document 1"}}'
# Search vectors
curl -X POST http://localhost:8082/api/v1/stores/my_store/search \
-H "Authorization: Bearer YOUR_TOKEN" \
-H "Content-Type: application/json" \
-d '{"vector": [1.1, 2.1, 3.1], "k": 5}'OAuth 2.1身份验证
默认凭据
系统在启动时创建默认的管理客户端:
- 客户端ID:
admin - 客户端密钥:
admin_secret_2024 - 范围:
read,write,admin
获取访问令牌
使用Erlang API:
{ok, TokenResponse} = erlvectordb:get_oauth_token(
>,
>,
[>, >]
).
AccessToken = maps:get(access_token, TokenResponse).使用HTTP API:
curl -X POST http://localhost:8081/oauth/token \
-H "Content-Type: application/x-www-form-urlencoded" \
-d "grant_type=client_credentials&client_id=admin&client_secret=admin_secret_2024&scope=read write"注册新客户
erlvectordb:register_oauth_client(>, >, #{
scopes => [>, >],
grant_types => [>, >]
}).基于范围的权限
read:搜索向量,获取统计信息write:插入/删除向量,同步存储admin:备份/还原操作、客户端管理
AI集成
Gemini AI集成
ErlVectorDB包括使用谷歌Gemini AI进行语义搜索和智能文档分析的AI功能:
# Setup Gemini integration
export GEMINI_API_KEY='your-gemini-api-key'
pip install google-generativeai
bash examples/setup_gemini_demo.sh
# Run AI-enhanced demo
python examples/gemini_mcp_client.pyAI功能:
- 语义嵌入:使用Gemini AI生成向量嵌入
- 智能文档分析:自动分类和元数据提取
- 智能搜索:带有AI解释的自然语言查询
- 内容理解:情感分析和主题提取
示例用法:
from examples.gemini_mcp_client import GeminiMCPClient
client = GeminiMCPClient(gemini_api_key="your-key")
client.connect_to_vectordb()
# AI-powered document insertion with automatic analysis
result = client.smart_insert('ai_store', 'healthcare_doc',
'AI is revolutionizing medical diagnosis and treatment...')
# Semantic search with natural language
search = client.smart_search('medical AI applications', 'ai_store')
print(search['explanation']) # AI explains why results are relevantGemini CLI配置
要将ErlVectorDB与Gemini CLI一起使用,您需要将其配置为MCP服务器。
步骤1:获取Gemini API密钥
从以下位置获取API密钥:https://makersuite.google.com/app/apikey
步骤2:配置Gemini CLI
创建或编辑Gemini CLI配置文件:
位置: ~/.config/gemini-cli/config.json (Linux/macOS)或 %USERPROFILE%\.config\gemini-cli\config.json (Windows)
{
"mcpServers": {
"erlvectordb": {
"command": "python",
"args": ["/path/to/Erlvectordb/examples/gemini_mcp_server.py"],
"env": {
"ERLVECTORDB_HOST": "localhost",
"ERLVECTORDB_PORT": "8080",
"ERLVECTORDB_OAUTH_HOST": "localhost",
"ERLVECTORDB_OAUTH_PORT": "8081",
"ERLVECTORDB_CLIENT_ID": "admin",
"ERLVECTORDB_CLIENT_SECRET": "admin_secret_2024"
}
}
}
}重要:替换 /path/to/Erlvectordb 使用ErlVectorDB安装的实际路径。
有关详细的配置选项,请参阅 Gemini MCP服务器配置指南.
步骤3:启动ErlVectorDB
# Using the local starter script (recommended)
./start-local.sh
# Or manually
rebar3 shell --eval "application:ensure_all_started(erlvectordb)"步骤4:与Gemini CLI一起使用
# Start Gemini CLI
gemini-cli
# The ErlVectorDB MCP server will be automatically available
# You can now use natural language to interact with your vector database:
> "Store this document about machine learning in the database"
> "Find similar documents about AI"
> "What documents do I have about healthcare?"Gemini CLI连接故障排除
如果出现“连接关闭”错误:
- 验证ErlVectorDB是否正在运行:
./check-status.sh- 直接测试MCP服务器:
echo '{"jsonrpc":"2.0","method":"initialize","id":1,"params":{}}' | \
python examples/gemini_mcp_server.py- 检查日志:
MCP服务器登录到stderr,因此您可以在Gemini CLI输出中看到错误。
- 验证配置中的路径:
确保路径 gemini_mcp_server.py 是绝对和正确的。
- 测试OAuth连接:
curl -X POST http://localhost:8081/oauth/token \
-d "grant_type=client_credentials&client_id=admin&client_secret=admin_secret_2024&scope=read write"- 检查Python依赖关系:
python3 -c "import requests; print('requests OK')"- 启用调试日志记录:
export ERLVECTORDB_LOG_LEVEL=DEBUG
python examples/gemini_mcp_server.pyGemini MCP服务器常见问题
问题:“连接超时”
- 症状:服务器无法连接到ErlVectorDB
- 解决方案:增加环境变量中的超时时间:
export ERLVECTORDB_SOCKET_TIMEOUT=60问题:“身份验证错误”
- 症状:OAuth令牌获取失败
- 解决方案:验证OAuth服务器是否正在运行以及凭据是否正确:
# Check OAuth server
curl http://localhost:8081/oauth/token \
-d "grant_type=client_credentials&client_id=admin&client_secret=admin_secret_2024&scope=read write"问题:“解析错误”
- 症状:JSON解析失败
- 解决方案:确保请求格式正确,为带换行符的单行JSON
问题:“远程主机关闭连接”
- 症状:ErlVectorDB意外关闭连接
- 解决方案:检查ErlVectorDB日志是否有错误,验证网络连接
问题:“消息不完整”
- 症状:大反应被截断
- 解决方案:增加缓冲区大小:
export ERLVECTORDB_BUFFER_SIZE=16384Gemini MCP服务器配置选项
这 gemini_mcp_server.py 脚本支持以下环境变量:
连接设置:
ERLVECTORDB_HOST-MCP服务器主机(默认:localhost)ERLVECTORDB_PORT-MCP服务器端口(默认值:8080)ERLVECTORDB_OAUTH_HOST-OAuth服务器主机(默认:localhost)ERLVECTORDB_OAUTH_PORT-OAuth服务器端口(默认:8081)
身份验证:
ERLVECTORDB_CLIENT_ID-OAuth客户端ID(默认值:admin)ERLVECTORDB_CLIENT_SECRET-OAuth客户端密钥(默认:admin_secret_2024)
插座配置:
ERLVECTORDB_SOCKET_TIMEOUT-套接字超时(秒)(默认值:30)ERLVECTORDB_BUFFER_SIZE-套接字缓冲区大小(字节)(默认值:8192)
重新连接设置:
ERLVECTORDB_MAX_RECONNECT_ATTEMPTS-最大重新连接尝试次数(默认值:3)ERLVECTORDB_RECONNECT_DELAY-初始重新连接延迟(秒)(默认值:1.0)
OAuth重试设置:
ERLVECTORDB_OAUTH_MAX_RETRIES-OAuth令牌请求重试次数上限(默认值:3)ERLVECTORDB_OAUTH_INITIAL_BACKOFF-初始退避延迟(秒)(默认值:1.0)ERLVECTORDB_OAUTH_MAX_BACKOFF-最大退避延迟(秒)(默认值:30.0)ERLVECTORDB_OAUTH_BACKOFF_MULTIPLIER-回退倍数(默认值:2.0)
登录中:
ERLVECTORDB_LOG_LEVEL-日志级别:调试、信息、警告、错误、严重(默认值:信息)
配置示例:
# High-latency network configuration
export ERLVECTORDB_SOCKET_TIMEOUT=60
export ERLVECTORDB_MAX_RECONNECT_ATTEMPTS=5
export ERLVECTORDB_RECONNECT_DELAY=2.0
# Large message handling
export ERLVECTORDB_BUFFER_SIZE=32768
# Debug mode
export ERLVECTORDB_LOG_LEVEL=DEBUG
# Run server
python examples/gemini_mcp_server.pyGemini MCP服务器使用示例
Gemini CLI的基本用法:
配置后,您可以在Gemini CLI中使用自然语言:
# Start Gemini CLI
gemini-cli
# Example interactions:
> "Create a vector store called 'documents'"
> "Insert a vector [1.0, 2.0, 3.0] with id 'doc1' into the documents store"
> "Search for vectors similar to [1.1, 2.1, 3.1] in the documents store"
> "List all available tools"直接测试(无Gemini CLI):
直接通过stdin/stdout测试MCP服务器:
# Test initialize
echo '{"jsonrpc":"2.0","method":"initialize","id":1,"params":{"protocolVersion":"2024-11-05","capabilities":{"tools":{}}}}' | \
python examples/gemini_mcp_server.py
# Test tools/list
echo '{"jsonrpc":"2.0","method":"tools/list","id":2,"params":{}}' | \
python examples/gemini_mcp_server.py
# Test tools/call
echo '{"jsonrpc":"2.0","method":"tools/call","id":3,"params":{"name":"create_store","arguments":{"name":"test_store"}}}' | \
python examples/gemini_mcp_server.py程序化使用:
在您自己的Python代码中使用服务器组件:
from examples.gemini_mcp_server import MCPServer, ServerConfig
# Create configuration
config = ServerConfig.from_environment()
# Or customize configuration
config = ServerConfig(
erlvectordb_host='localhost',
erlvectordb_port=8080,
oauth_host='localhost',
oauth_port=8081,
client_id='admin',
client_secret='admin_secret_2024',
socket_timeout=60,
log_level='DEBUG'
)
# Validate configuration
config.validate()
# Create and run server
server = MCPServer(config)
exit_code = server.run()组件级别使用:
使用单个组件进行自定义集成:
from examples.gemini_mcp_server import (
SocketHandler, OAuthManager, RequestRouter,
StdioHandler, ServerConfig
)
# Create configuration
config = ServerConfig.from_environment()
# OAuth token management
oauth_manager = OAuthManager(config)
token = oauth_manager.get_token()
print(f"Access token: {token}")
# Socket communication
socket_handler = SocketHandler(
host=config.erlvectordb_host,
port=config.erlvectordb_port
)
socket_handler.connect()
# Send request
request = {
'jsonrpc': '2.0',
'method': 'tools/list',
'params': {},
'id': 1,
'auth': oauth_manager.get_auth_dict()
}
socket_handler.send_message(request)
response = socket_handler.receive_message()
print(f"Response: {response}")
# Cleanup
socket_handler.close()替代方案:直接Python脚本
如果你更喜欢直接使用Python客户端进行演示:
# Set environment variables
export GEMINI_API_KEY='your-gemini-api-key'
export ERLVECTORDB_HOST='localhost'
export ERLVECTORDB_PORT='8080'
export ERLVECTORDB_OAUTH_HOST='localhost'
export ERLVECTORDB_OAUTH_PORT='8081'
export ERLVECTORDB_CLIENT_ID='admin'
export ERLVECTORDB_CLIENT_SECRET='admin_secret_2024'
# Run the demo client (not for Gemini CLI)
python examples/gemini_mcp_client.pyClaude桌面集成:
{
"mcpServers": {
"erlvectordb-ai": {
"command": "python",
"args": ["examples/gemini_mcp_client.py"],
"env": {
"GEMINI_API_KEY": "${GEMINI_API_KEY}",
"ERLVECTORDB_HOST": "localhost",
"ERLVECTORDB_PORT": "8080"
}
}
}
}MCP集成
该数据库在端口8080(可配置)上公开了一个具有OAuth 2.1身份验证的MCP服务器。OAuth服务器在端口8081上运行。
📖 详细设置指南:参见 MCP设置指南 获取完整的配置和连接说明。
快速MCP设置
- 启动ErlVectorDB:
rebar3 shell- 测试连接:
# Get OAuth token
curl -X POST http://localhost:8081/oauth/token \
-d "grant_type=client_credentials&client_id=admin&client_secret=admin_secret_2024&scope=read write"
# Test MCP connection
telnet localhost 8080- 运行示例客户端:
node examples/mcp_client.js
# or
python examples/mcp_client.py
# or AI-enhanced with Gemini
export GEMINI_API_KEY='your-key'
python examples/gemini_mcp_client.pyMCP服务器配置
基本配置
% In config/sys.config
[
{erlvectordb, [
{mcp_port, 8080}, % MCP server port
{oauth_enabled, true}, % Enable OAuth authentication
{oauth_port, 8081}, % OAuth server port
% ... other settings
]}
].禁用身份验证(仅限开发)
{oauth_enabled, false} % Disables OAuth - NOT recommended for productionMCP客户端配置
适用于克劳德桌面
添加到您的Claude Desktop配置文件中:
macOS: ~/Library/Application Support/Claude/claude_desktop_config.json 视窗: %APPDATA%\Claude\claude_desktop_config.json
{
"mcpServers": {
"erlvectordb": {
"command": "node",
"args": ["/path/to/mcp-client.js"],
"env": {
"ERLVECTORDB_HOST": "localhost",
"ERLVECTORDB_PORT": "8080",
"ERLVECTORDB_OAUTH_HOST": "localhost",
"ERLVECTORDB_OAUTH_PORT": "8081",
"ERLVECTORDB_CLIENT_ID": "admin",
"ERLVECTORDB_CLIENT_SECRET": "admin_secret_2024"
}
}
}
}对于其他MCP客户端
配置您的MCP客户端以连接到:
- MCP端点:
tcp://localhost:8080 - OAuth端点:
http://localhost:8081/oauth/token
MCP客户端实现
以下是一个完整的Node.js MCP客户端示例:
// mcp-client.js
const net = require('net');
const https = require('https');
class ErlVectorDBClient {
constructor(options = {}) {
this.host = options.host || 'localhost';
this.port = options.port || 8080;
this.oauthHost = options.oauthHost || 'localhost';
this.oauthPort = options.oauthPort || 8081;
this.clientId = options.clientId || 'admin';
this.clientSecret = options.clientSecret || 'admin_secret_2024';
this.accessToken = null;
this.socket = null;
}
async getAccessToken() {
const postData = new URLSearchParams({
grant_type: 'client_credentials',
client_id: this.clientId,
client_secret: this.clientSecret,
scope: 'read write admin'
}).toString();
const options = {
hostname: this.oauthHost,
port: this.oauthPort,
path: '/oauth/token',
method: 'POST',
headers: {
'Content-Type': 'application/x-www-form-urlencoded',
'Content-Length': Buffer.byteLength(postData)
}
};
return new Promise((resolve, reject) => {
const req = http.request(options, (res) => {
let data = '';
res.on('data', (chunk) => data += chunk);
res.on('end', () => {
try {
const response = JSON.parse(data);
if (response.access_token) {
this.accessToken = response.access_token;
resolve(response.access_token);
} else {
reject(new Error('No access token received'));
}
} catch (error) {
reject(error);
}
});
});
req.on('error', reject);
req.write(postData);
req.end();
});
}
async connect() {
if (!this.accessToken) {
await this.getAccessToken();
}
return new Promise((resolve, reject) => {
this.socket = net.createConnection(this.port, this.host, () => {
console.log('Connected to ErlVectorDB MCP server');
resolve();
});
this.socket.on('error', reject);
this.socket.on('data', (data) => {
try {
const response = JSON.parse(data.toString());
this.handleResponse(response);
} catch (error) {
console.error('Failed to parse response:', error);
}
});
});
}
async sendRequest(method, params = {}, id = 1) {
const request = {
jsonrpc: '2.0',
method: method,
params: params,
id: id,
auth: {
type: 'bearer',
token: this.accessToken
}
};
return new Promise((resolve, reject) => {
this.responseHandlers = this.responseHandlers || {};
this.responseHandlers[id] = { resolve, reject };
const requestData = JSON.stringify(request);
this.socket.write(requestData);
});
}
handleResponse(response) {
if (response.id && this.responseHandlers[response.id]) {
const handler = this.responseHandlers[response.id];
delete this.responseHandlers[response.id];
if (response.error) {
handler.reject(new Error(response.error.message));
} else {
handler.resolve(response.result);
}
}
}
// MCP Protocol Methods
async initialize() {
return this.sendRequest('initialize', {
protocolVersion: '2024-11-05',
capabilities: {
tools: {}
},
clientInfo: {
name: 'erlvectordb-client',
version: '1.0.0'
}
});
}
async listTools() {
return this.sendRequest('tools/list');
}
async callTool(name, arguments) {
return this.sendRequest('tools/call', {
name: name,
arguments: arguments
});
}
// Vector Database Operations
async createStore(name) {
return this.callTool('create_store', { name: name });
}
async insertVector(store, id, vector, metadata = {}) {
return this.callTool('insert_vector', {
store: store,
id: id,
vector: vector,
metadata: metadata
});
}
async searchVectors(store, vector, k = 10) {
return this.callTool('search_vectors', {
store: store,
vector: vector,
k: k
});
}
async syncStore(store) {
return this.callTool('sync_store', { store: store });
}
async backupStore(store, backupName) {
return this.callTool('backup_store', {
store: store,
backup_name: backupName
});
}
async listBackups() {
return this.callTool('list_backups', {});
}
disconnect() {
if (this.socket) {
this.socket.end();
this.socket = null;
}
}
}
// Usage Example
async function main() {
const client = new ErlVectorDBClient({
host: process.env.ERLVECTORDB_HOST || 'localhost',
port: parseInt(process.env.ERLVECTORDB_PORT) || 8080,
oauthHost: process.env.ERLVECTORDB_OAUTH_HOST || 'localhost',
oauthPort: parseInt(process.env.ERLVECTORDB_OAUTH_PORT) || 8081,
clientId: process.env.ERLVECTORDB_CLIENT_ID || 'admin',
clientSecret: process.env.ERLVECTORDB_CLIENT_SECRET || 'admin_secret_2024'
});
try {
await client.connect();
await client.initialize();
const tools = await client.listTools();
console.log('Available tools:', tools);
// Create a store
await client.createStore('test_store');
// Insert a vector
await client.insertVector('test_store', 'doc1', [1.0, 2.0, 3.0], {
title: 'Test Document'
});
// Search for similar vectors
const results = await client.searchVectors('test_store', [1.1, 2.1, 3.1], 5);
console.log('Search results:', results);
} catch (error) {
console.error('Error:', error);
} finally {
client.disconnect();
}
}
if (require.main === module) {
main();
}
module.exports = ErlVectorDBClient;Python MCP客户端
# mcp_client.py
import json
import socket
import requests
from typing import Dict, List, Any, Optional
class ErlVectorDBClient:
def __init__(self, host='localhost', port=8080, oauth_host='localhost',
oauth_port=8081, client_id='admin', client_secret='admin_secret_2024'):
self.host = host
self.port = port
self.oauth_host = oauth_host
self.oauth_port = oauth_port
self.client_id = client_id
self.client_secret = client_secret
self.access_token = None
self.socket = None
self.request_id = 1
def get_access_token(self) -> str:
"""Get OAuth access token"""
url = f'http://{self.oauth_host}:{self.oauth_port}/oauth/token'
data = {
'grant_type': 'client_credentials',
'client_id': self.client_id,
'client_secret': self.client_secret,
'scope': 'read write admin'
}
response = requests.post(url, data=data)
response.raise_for_status()
token_data = response.json()
self.access_token = token_data['access_token']
return self.access_token
def connect(self):
"""Connect to MCP server"""
if not self.access_token:
self.get_access_token()
self.socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
self.socket.connect((self.host, self.port))
def send_request(self, method: str, params: Dict = None, request_id: int = None) -> Dict:
"""Send MCP request"""
if request_id is None:
request_id = self.request_id
self.request_id += 1
request = {
'jsonrpc': '2.0',
'method': method,
'params': params or {},
'id': request_id,
'auth': {
'type': 'bearer',
'token': self.access_token
}
}
request_data = json.dumps(request).encode('utf-8')
self.socket.send(request_data)
# Receive response
response_data = self.socket.recv(4096)
response = json.loads(response_data.decode('utf-8'))
if 'error' in response:
raise Exception(f"MCP Error: {response['error']['message']}")
return response.get('result', {})
def initialize(self) -> Dict:
"""Initialize MCP connection"""
return self.send_request('initialize', {
'protocolVersion': '2024-11-05',
'capabilities': {'tools': {}},
'clientInfo': {
'name': 'erlvectordb-python-client',
'version': '1.0.0'
}
})
def list_tools(self) -> Dict:
"""List available tools"""
return self.send_request('tools/list')
def call_tool(self, name: str, arguments: Dict) -> Dict:
"""Call a specific tool"""
return self.send_request('tools/call', {
'name': name,
'arguments': arguments
})
# Vector Database Operations
def create_store(self, name: str) -> Dict:
return self.call_tool('create_store', {'name': name})
def insert_vector(self, store: str, vector_id: str, vector: List[float],
metadata: Dict = None) -> Dict:
return self.call_tool('insert_vector', {
'store': store,
'id': vector_id,
'vector': vector,
'metadata': metadata or {}
})
def search_vectors(self, store: str, vector: List[float], k: int = 10) -> Dict:
return self.call_tool('search_vectors', {
'store': store,
'vector': vector,
'k': k
})
def sync_store(self, store: str) -> Dict:
return self.call_tool('sync_store', {'store': store})
def backup_store(self, store: str, backup_name: str) -> Dict:
return self.call_tool('backup_store', {
'store': store,
'backup_name': backup_name
})
def list_backups(self) -> Dict:
return self.call_tool('list_backups', {})
def disconnect(self):
"""Disconnect from server"""
if self.socket:
self.socket.close()
self.socket = None
# Usage Example
if __name__ == '__main__':
client = ErlVectorDBClient()
try:
client.connect()
client.initialize()
# List available tools
tools = client.list_tools()
print('Available tools:', tools)
# Create a store
client.create_store('python_test_store')
# Insert vectors
client.insert_vector('python_test_store', 'doc1', [1.0, 2.0, 3.0],
{'title': 'Python Test Document'})
# Search
results = client.search_vectors('python_test_store', [1.1, 2.1, 3.1], 5)
print('Search results:', results)
except Exception as e:
print(f'Error: {e}')
finally:
client.disconnect()可用工具(按范围):
读取范围:
search_vectors:搜索相似向量
写入范围:
create_store:创建新的矢量存储insert_vector:插入包含元数据的向量sync_store:将存储同步到持久存储
管理范围:
backup_store:创建存储的备份restore_store:从备份还原存储list_backups:列出所有可用备份
使用OAuth的MCP客户端示例
{
"jsonrpc": "2.0",
"method": "tools/call",
"params": {
"name": "insert_vector",
"arguments": {
"store": "my_store",
"id": "doc1",
"vector": [1.0, 2.0, 3.0],
"metadata": {"title": "Document 1"}
}
},
"auth": {
"type": "bearer",
"token": "your_access_token_here"
},
"id": 1
}MCP连接测试
测试OAuth连接
# Test OAuth token endpoint
curl -X POST http://localhost:8081/oauth/token \
-H "Content-Type: application/x-www-form-urlencoded" \
-d "grant_type=client_credentials&client_id=admin&client_secret=admin_secret_2024&scope=read write admin"测试MCP连接
# Test MCP server connectivity
telnet localhost 8080
# Send initialize request (after connecting)
{"jsonrpc":"2.0","method":"initialize","params":{"protocolVersion":"2024-11-05","capabilities":{"tools":{}}},"id":1}MCP集成模式
流媒体响应
对于大型结果集,实现流式传输:
// Handle streaming responses
client.socket.on('data', (chunk) => {
const lines = chunk.toString().split('\n');
lines.forEach(line => {
if (line.trim()) {
try {
const response = JSON.parse(line);
handleStreamingResponse(response);
} catch (e) {
// Handle partial JSON
buffer += line;
}
}
});
});连接池
对于高通量应用:
class ErlVectorDBPool {
constructor(options, poolSize = 5) {
this.options = options;
this.poolSize = poolSize;
this.connections = [];
this.available = [];
}
async getConnection() {
if (this.available.length > 0) {
return this.available.pop();
}
if (this.connections.length {
const checkAvailable = () => {
if (this.available.length > 0) {
resolve(this.available.pop());
} else {
setTimeout(checkAvailable, 10);
}
};
checkAvailable();
});
}
releaseConnection(client) {
this.available.push(client);
}
}错误处理和重试逻辑
class RobustErlVectorDBClient extends ErlVectorDBClient {
async sendRequestWithRetry(method, params, maxRetries = 3) {
for (let attempt = 1; attempt
setTimeout(resolve, Math.pow(2, attempt) * 1000)
);
}
}
}
}MCP故障排除
常见问题
1.身份验证错误
Error: Authentication required- 验证OAuth服务器是否在端口8081上运行
- 检查客户端凭据是否正确
- 确保令牌未过期(默认值:1小时)
2.连接被拒绝
Error: ECONNREFUSED- 验证ErlVectorDB是否正在运行
- 检查MCP端口配置(默认值:8080)
- 确保防火墙允许连接
3.无效的工具调用
Error: Insufficient permissions- 检查OAuth作用域是否符合工具要求
- 验证客户端是否具有必要的权限
- 审查基于范围的访问控制
4.JSON解析错误
Error: Unexpected token in JSON- 确保JSON格式正确
- 检查尾随逗号或语法错误
- 验证UTF-8编码
调试模式
在ErlVectorDB中启用调试日志记录:
% In config/sys.config
{kernel, [
{logger_level, debug},
{logger, [
{handler, default, logger_std_h,
#{config => #{type => standard_io},
formatter => {logger_formatter, #{}}}}
]}
]}健康检查端点
测试服务器运行状况:
# Check if servers are responding
curl -f http://localhost:8081/oauth/client_info \
-H "Authorization: Bearer YOUR_TOKEN" || echo "OAuth server down"
telnet localhost 8080 || echo "MCP server down"MCP性能优化
批量操作
// Batch insert multiple vectors
async function batchInsert(client, store, vectors) {
const promises = vectors.map(({id, vector, metadata}) =>
client.insertVector(store, id, vector, metadata)
);
return Promise.all(promises);
}连接保持活跃
// Implement heartbeat to keep connection alive
setInterval(() => {
if (client.socket && !client.socket.destroyed) {
client.listTools().catch(err => {
console.log('Heartbeat failed, reconnecting...');
client.connect();
});
}
}, 30000); // 30 seconds大向量压缩
// Use compression for large vectors
async function insertLargeVector(client, store, id, vector, metadata) {
if (vector.length > 1000) {
// Let server handle compression
return client.callTool('insert_compressed_vector', {
store, id, vector, metadata
});
} else {
return client.insertVector(store, id, vector, metadata);
}
}配置
编辑 config/sys.config:
[
{erlvectordb, [
% Network configuration
{mcp_port, 8080},
{oauth_port, 8081},
{rest_api_port, 8082},
{rest_api_enabled, true},
% Authentication
{oauth_enabled, true},
{create_default_client, true},
{default_client_id, >},
{default_client_secret, >},
{token_lifetime, 3600000}, % 1 hour
{refresh_token_lifetime, 86400000}, % 24 hours
% Clustering
{cluster_enabled, false},
{node_name, 'erlvectordb@localhost'},
{cluster_cookie, erlvectordb_cluster},
{replication_factor, 2},
{heartbeat_interval, 5000},
% Compression
{compression_enabled, true},
{compression_algorithm, quantization_8bit},
% Storage
{persistence_enabled, true},
{persistence_dir, "data"},
{backup_dir, "backups"},
{sync_interval, 30000}
]}
].测试
rebar3 ct性能特征
- 并发操作:每个矢量存储都在自己的进程中运行
- 并行搜索:搜索操作可以跨存储和节点同时运行
- 内存效率高:利用Erlang的写时复制语义+压缩
- 容错:个别商店的崩溃不会影响其他商店
- 水平扩展:使用集群节点进行线性缩放
- 压缩优势:通过可配置的精度权衡,减少50-90%的存储空间
- 网络优化:压缩复制减少了带宽使用
持久性和备份
自动持久化
所有向量操作都使用DETS(disk Erlang Term Storage)自动持久化到磁盘,ETS用于快速内存访问:
% Data is automatically saved, but you can force sync
erlvectordb:sync(my_store).备份操作
% Create a backup
{ok, BackupInfo} = erlvectordb:backup_store(my_store, "production_backup"),
BackupPath = BackupInfo#backup_info.file_path.
% List all backups
{ok, Backups} = erlvectordb:list_backups().
% Restore from backup to new store
{ok, Result} = erlvectordb:restore_store(BackupPath, new_store_name).导出/导入
% Export to JSON format
{ok, _} = erlvectordb:export_store(my_store, "export.json").
% Import from JSON
{ok, _} = erlvectordb:import_store("export.json", imported_store).高级功能
自定义距离度量
% Use vector_utils for custom operations
Similarity = vector_utils:cosine_similarity(Vector1, Vector2),
Distance = vector_utils:euclidean_distance(Vector1, Vector2).索引管理
% Create advanced indexes (future feature)
vector_index_manager:create_index(my_index, #{type => hnsw, m => 16}).许可证
此项目根据GNU通用公共许可证v3.0获得许可-请参阅 许可证 文件以获取详细信息。
贡献
- 克隆该仓库
- 创建要素分支
- 添加新功能的测试
- 确保所有测试通过:
rebar3 ct - 提交拉取请求
路线图
- \[x\] 使用DETS/ETS进行持久存储
- \[x\] 备份和恢复系统
- \[x\] JSON导出/导入功能
- \[x\] 基于作用域权限的OAuth 2.1身份验证
- \[x\] 具有自动复制功能的分布式群集
- \[x\] REST API与MCP协议
- \[x\] 采用多种算法的矢量压缩
- \[\]用于近似最近邻搜索的HNSW索引
- \[\]批量操作以提高吞吐量
- \[\]矢量量化优化
- \[\]普罗米修斯指标集成
- \[\]自动备份计划
- \[\]高级压缩算法(LSH、随机投影)
- \[\]JWT令牌支持
- \[\]公共客户端的PKCE流
- \[\]GraphQL API
- \[\]WebSocket实时更新
架构优势
此实现利用了Erlang/OTP的独特优势:
- 故障隔离:每个矢量存储在其自己的进程中都是隔离的
- 热代码重新加载:在不停止数据库的情况下更新算法
- 大规模并发:处理数千个并发操作
- 分布:跨多个节点的本机群集,具有自动故障转移功能
- 监督:故障组件的自动重启
- 模式匹配:高效的消息路由和数据处理
- 永久存储:DETS提供耐久性和ETS性能
- 备份系统:具有JSON互操作性的内置备份/还原
- OAuth 2.1安全:具有基于范围的访问控制的行业标准身份验证
- 压缩:智能矢量压缩可将存储空间减少50-90%
- REST集成:与MCP协议一起用于web应用程序的完整REST API
