Canner / WrenAI 全面解析:新一代数据智能平台实战指南
在日常数据开发工作中,你是否遇到过这样的困境:数据源分散在不同系统,SQL编写效率低下,业务人员难以自主分析数据?传统的数据平台往往需要专业的数据工程师进行复杂配置,业务团队的数据需求响应缓慢。本文将深入解析Canner和WrenAI这两个新一代数据智能平台,从核心概念到实战应用,带你掌握现代化数据平台的建设思路。
1. 背景与核心概念
1.1 什么是Canner和WrenAI
Canner和WrenAI都是面向现代数据栈的智能数据平台,旨在简化数据工程流程,提升数据分析效率。Canner主要专注于数据虚拟化领域,提供统一的数据访问层,让用户能够通过单一接口访问分布在多个数据源中的数据。而WrenAI则更侧重于AI驱动的数据分析和自然语言查询,让非技术用户也能轻松进行复杂的数据探索。
这两个平台都体现了现代数据平台的发展趋势:降低技术门槛、提升自动化程度、增强用户体验。与传统的数据平台相比,它们更加注重业务人员的需求,通过智能化的方式减少中间环节,实现数据价值的快速释放。
1.2 核心价值与解决的问题
Canner和WrenAI主要解决企业在数据应用过程中面临的几个核心痛点。首先是数据孤岛问题,企业中的数据往往分散在数据库、数据仓库、数据湖等多个系统中,难以统一管理和使用。其次是技术门槛问题,传统的数据分析需要专业的SQL技能,限制了业务人员的自主分析能力。最后是效率问题,从数据需求提出到最终结果产出,往往需要经历漫长的开发周期。
通过数据虚拟化和AI技术,这两个平台能够实现数据的统一访问、智能查询优化和自然语言交互,显著提升数据应用的效率和覆盖面。业务人员可以直接用自然语言提出数据需求,系统自动生成相应的查询语句并返回结果,大大缩短了数据价值实现的路径。
1.3 典型应用场景
在实际业务中,Canner和WrenAI适用于多种场景。对于需要快速响应业务变化的企业,可以通过这些平台建立敏捷的数据服务体系。比如在零售行业,营销团队需要实时查看销售数据、用户行为数据,传统方式需要向数据团队提需求、排队等待开发,而现在业务人员可以直接用自然语言查询所需数据。
在金融风控场景中,风险分析师需要综合多个数据源的信息进行决策,Canner的数据虚拟化能力可以统一访问交易数据、用户画像数据、外部风险数据等,而WrenAI的智能分析能力可以帮助发现潜在的风险模式。此外,在制造业、医疗健康、教育等多个行业,这些平台都能发挥重要作用。
2. 技术架构与核心特性
2.1 Canner的技术架构
Canner的核心架构基于数据虚拟化技术,采用分层设计理念。最底层是数据连接层,支持各种类型的数据源接入,包括关系型数据库(MySQL、PostgreSQL等)、数据仓库(Snowflake、BigQuery等)、数据湖(S3、ADLS等)以及API数据源。中间层是查询优化层,负责将用户提交的SQL查询分解为针对不同数据源的具体查询计划,并进行性能优化。
最上层是统一服务层,提供标准的SQL接口和API接口,应用程序可以通过这些接口透明地访问底层所有数据源。Canner还内置了数据缓存、查询重写、下推优化等高级功能,确保查询性能达到最优。这种架构使得用户无需关心数据的物理存储位置,只需关注业务逻辑本身。
2.2 WrenAI的智能特性
WrenAI的核心竞争力在于其AI驱动能力。首先是最突出的自然语言转SQL功能,用户可以用日常语言描述数据需求,系统通过大语言模型理解用户意图,自动生成准确的SQL查询语句。其次是智能数据发现功能,系统可以自动分析数据之间的关系,推荐相关的数据分析和可视化方案。
另一个重要特性是查询优化建议,WrenAI可以分析查询模式,提出性能优化建议,甚至自动重写低效的查询语句。在数据治理方面,WrenAI能够自动识别数据质量问题,发现数据异常,并提供数据血缘分析等高级功能。这些智能特性使得数据平台从被动的工具转变为主动的智能助手。
2.3 平台集成能力
两个平台都具备强大的集成能力,可以与现代数据栈中的各种工具无缝对接。在数据源支持方面,除了常见的关系数据库和数据仓库,还支持流式数据源、NoSQL数据库、云存储服务等。在BI工具集成方面,支持Tableau、Power BI、Superset等主流可视化工具。
对于开发团队,平台提供完整的API接口和SDK支持,可以方便地嵌入到现有应用中。还支持与调度工具(如Airflow)、监控工具(如Prometheus)、身份认证系统(如OAuth)的集成。这种开放的架构设计确保了平台可以灵活地融入企业现有的技术生态。
3. 环境准备与安装部署
3.1 系统要求与依赖环境
在部署Canner或WrenAI之前,需要确保环境满足基本要求。操作系统方面,支持Linux(CentOS 7+、Ubuntu 18.04+)和Windows Server 2016+,推荐使用Linux系统以获得更好的性能。硬件配置建议至少4核CPU、8GB内存、100GB存储空间,具体规模取决于数据量和并发用户数。
软件依赖包括Docker 20.10+、Docker Compose 1.29+,如果选择Kubernetes部署则需要k8s 1.20+版本。网络方面需要确保服务器可以访问目标数据源,如果数据源在云端,还需要配置相应的网络权限。对于生产环境部署,建议使用负载均衡、高可用配置,确保服务稳定性。
3.2 Docker快速部署
对于测试和开发环境,推荐使用Docker Compose进行快速部署。首先创建部署目录和配置文件:
# 创建项目目录 mkdir canner-deployment && cd canner-deployment # 创建docker-compose.yml文件 cat > docker-compose.yml << 'EOF' version: '3.8' services: canner-server: image: canner/canner:latest ports: - "8080:8080" environment: - CANNER_DB_TYPE=postgresql - CANNER_DB_HOST=postgres - CANNER_DB_PORT=5432 - CANNER_DB_NAME=canner - CANNER_DB_USER=canner_user - CANNER_DB_PASSWORD=your_password depends_on: - postgres volumes: - ./config:/app/config postgres: image: postgres:13 environment: - POSTGRES_DB=canner - POSTGRES_USER=canner_user - POSTGRES_PASSWORD=your_password volumes: - postgres_data:/var/lib/postgresql/data volumes: postgres_data: EOF启动服务:
# 启动服务 docker-compose up -d # 检查服务状态 docker-compose ps # 查看日志 docker-compose logs canner-server3.3 Kubernetes生产部署
对于生产环境,建议使用Kubernetes进行部署。首先创建命名空间和配置文件:
# canner-namespace.yaml apiVersion: v1 kind: Namespace metadata: name: canner# canner-deployment.yaml apiVersion: apps/v1 kind: Deployment metadata: name: canner-server namespace: canner spec: replicas: 3 selector: matchLabels: app: canner-server template: metadata: labels: app: canner-server spec: containers: - name: canner-server image: canner/canner:latest ports: - containerPort: 8080 env: - name: CANNER_DB_TYPE value: "postgresql" - name: CANNER_DB_HOST value: "postgres-service" - name: CANNER_DB_PORT value: "5432" resources: requests: memory: "512Mi" cpu: "250m" limits: memory: "2Gi" cpu: "1000m" livenessProbe: httpGet: path: /health port: 8080 initialDelaySeconds: 30 periodSeconds: 10 --- apiVersion: v1 kind: Service metadata: name: canner-service namespace: canner spec: selector: app: canner-server ports: - port: 80 targetPort: 8080 type: LoadBalancer部署命令:
# 创建命名空间 kubectl apply -f canner-namespace.yaml # 部署服务 kubectl apply -f canner-deployment.yaml # 检查部署状态 kubectl get pods -n canner kubectl get services -n canner4. 数据源配置与管理
4.1 连接常见数据源
配置数据源是使用Canner/WrenAI的第一步,下面以MySQL和Snowflake为例演示配置过程。首先通过管理界面或API添加数据源:
// MySQL数据源配置示例 { "name": "production_mysql", "type": "mysql", "host": "mysql.example.com", "port": 3306, "database": "business_db", "username": "data_user", "password": "secure_password", "properties": { "useSSL": "true", "serverTimezone": "UTC", "connectTimeout": "30000" } } // Snowflake数据源配置示例 { "name": "analytics_warehouse", "type": "snowflake", "account": "company_account", "warehouse": "analytics_wh", "database": "analytics_db", "schema": "public", "username": "snowflake_user", "password": "secure_password", "role": "analyst_role" }对于WrenAI,还需要配置数据源的元数据信息,以便AI模型更好地理解数据结构:
# 数据源元数据配置 data_sources: - name: sales_mysql type: mysql description: "销售业务数据库,包含订单、客户、产品信息" tables: - name: orders description: "订单主表,记录所有销售订单" columns: - name: order_id description: "订单唯一标识" data_type: "bigint" - name: customer_id description: "客户ID" data_type: "int" - name: order_date description: "订单创建日期" data_type: "date" - name: customers description: "客户信息表"4.2 数据源权限管理
在企业环境中,数据安全至关重要。Canner/WrenAI提供细粒度的权限控制机制:
-- 创建角色和权限 CREATE ROLE business_analyst; GRANT USAGE ON DATABASE sales_db TO ROLE business_analyst; GRANT SELECT ON TABLE sales_db.orders TO ROLE business_analyst; GRANT SELECT ON TABLE sales_db.customers TO ROLE business_analyst; -- 创建用户并分配角色 CREATE USER alice IDENTIFIED BY 'password'; GRANT ROLE business_analyst TO USER alice; -- 列级权限控制 GRANT SELECT (order_id, order_date, amount) ON sales_db.orders TO ROLE business_analyst;通过API进行权限管理:
import requests import json # 配置API端点和管理员凭证 api_url = "https://canner.example.com/api/v1" admin_token = "your_admin_token" headers = { "Authorization": f"Bearer {admin_token}", "Content-Type": "application/json" } # 创建数据源访问策略 policy_data = { "name": "sales_data_policy", "description": "销售团队数据访问策略", "rules": [ { "data_sources": ["sales_mysql"], "tables": ["orders", "customers"], "allowed_operations": ["SELECT"], "row_filters": ["region = 'North'"], "column_masks": { "customers.phone": "partial(phone, 3, 4, '****')" } } ] } response = requests.post(f"{api_url}/policies", headers=headers, data=json.dumps(policy_data))4.3 数据源监控与健康检查
确保数据源连接的稳定性是平台可靠性的基础。配置健康检查机制:
# 健康检查配置 health_check: enabled: true interval: 300 # 5分钟检查一次 timeout: 30 # 30秒超时 thresholds: failure: 3 # 连续3次失败标记为不可用 checks: - name: mysql_connection type: jdbc query: "SELECT 1" expected_result: 1 - name: snowflake_warehouse type: snowflake query: "SELECT CURRENT_WAREHOUSE()" - name: api_response_time type: http url: "https://api.example.com/health" expected_status: 200 max_response_time: 1000 # 1秒内响应监控仪表板配置:
# 监控数据收集脚本 import time import psutil import requests from prometheus_client import start_http_server, Gauge, Counter # 定义监控指标 connection_errors = Counter('canner_connection_errors', 'Data source connection errors', ['data_source']) query_duration = Gauge('canner_query_duration_seconds', 'Query execution duration', ['data_source', 'status']) active_connections = Gauge('canner_active_connections', 'Active connections per data source', ['data_source']) def monitor_data_sources(): while True: for ds in data_sources: try: start_time = time.time() # 测试连接 result = test_connection(ds) duration = time.time() - start_time query_duration.labels(data_source=ds.name, status='success').set(duration) active_connections.labels(data_source=ds.name).set( get_active_connection_count(ds) ) except Exception as e: connection_errors.labels(data_source=ds.name).inc() query_duration.labels(data_source=ds.name, status='error').set(0) time.sleep(60) # 每分钟检查一次 if __name__ == '__main__': start_http_server(8000) monitor_data_sources()5. 查询优化与性能调优
5.1 查询执行计划分析
理解查询执行计划是性能优化的基础。Canner提供详细的执行计划分析功能:
-- 查看查询执行计划 EXPLAIN (FORMAT JSON) SELECT o.order_id, c.customer_name, SUM(oi.amount) as total_amount FROM orders o JOIN customers c ON o.customer_id = c.customer_id JOIN order_items oi ON o.order_id = oi.order_id WHERE o.order_date >= '2024-01-01' GROUP BY o.order_id, c.customer_name HAVING SUM(oi.amount) > 1000; -- 执行计划输出示例 { "plan": { "node_type": "Aggregate", "strategy": "Hashed", "plan_rows": 1000, "plan_width": 56, "actual_rows": 850, "actual_time": 15.234, "children": [ { "node_type": "Hash Join", "parent_relationship": "Outer", "join_type": "Inner", "plan_rows": 10000, "actual_rows": 12000, "actual_time": 8.765 } ] } }针对执行计划进行优化:
-- 优化前:全表扫描 SELECT * FROM orders WHERE YEAR(order_date) = 2024; -- 优化后:使用索引友好的条件 SELECT * FROM orders WHERE order_date >= '2024-01-01' AND order_date < '2025-01-01'; -- 添加合适的索引 CREATE INDEX idx_orders_date ON orders(order_date); CREATE INDEX idx_orders_customer_date ON orders(customer_id, order_date); -- 使用覆盖索引 CREATE INDEX idx_orders_covering ON orders(order_date, customer_id, amount);5.2 缓存策略配置
合理的缓存配置可以显著提升查询性能:
# 缓存配置示例 caching: enabled: true strategy: "adaptive" # 自适应缓存策略 # 查询结果缓存 query_result: enabled: true ttl: 3600 # 1小时 max_size: "10GB" eviction_policy: "LRU" # 元数据缓存 metadata: enabled: true ttl: 86400 # 24小时 # 执行计划缓存 plan_cache: enabled: true size: 1000 # 缓存1000个执行计划 # 自适应缓存规则 adaptive_rules: - min_execution_time: 1.0 # 执行时间超过1秒的查询 min_frequency: 5 # 最近被调用5次以上 cache_ttl: 7200 # 缓存2小时 - pattern: "SELECT.*FROM sales.*WHERE date = CURRENT_DATE" cache_ttl: 300 # 当前日期的查询缓存5分钟缓存监控和管理API:
import requests import json class CacheManager: def __init__(self, base_url, auth_token): self.base_url = base_url self.headers = { 'Authorization': f'Bearer {auth_token}', 'Content-Type': 'application/json' } def get_cache_stats(self): """获取缓存统计信息""" response = requests.get( f"{self.base_url}/api/v1/cache/stats", headers=self.headers ) return response.json() def clear_cache(self, cache_type=None, pattern=None): """清理缓存""" data = {} if cache_type: data['cache_type'] = cache_type if pattern: data['pattern'] = pattern response = requests.post( f"{self.base_url}/api/v1/cache/clear", headers=self.headers, data=json.dumps(data) ) return response.json() def preload_cache(self, queries): """预加载缓存""" data = {'queries': queries} response = requests.post( f"{self.base_url}/api/v1/cache/preload", headers=self.headers, data=json.dumps(data) ) return response.json() # 使用示例 cache_mgr = CacheManager('https://canner.example.com', 'your_token') stats = cache_mgr.get_cache_stats() print(f"缓存命中率: {stats['hit_rate']:.2%}")5.3 分布式查询优化
对于跨数据源的复杂查询,分布式查询优化至关重要:
-- 跨数据源查询示例 SELECT c.customer_name, o.order_date, p.product_name, SUM(oi.quantity) as total_quantity FROM mysql_sales.customers c JOIN snowflake_orders.orders o ON c.customer_id = o.customer_id JOIN bigquery_products.products p ON o.product_id = p.product_id JOIN redshift_items.order_items oi ON o.order_id = oi.order_id WHERE o.order_date BETWEEN '2024-01-01' AND '2024-03-31' GROUP BY c.customer_name, o.order_date, p.product_name HAVING SUM(oi.quantity) > 100; -- 查询优化策略 -- 1. 谓词下推:将过滤条件推送到数据源层执行 -- 2. 列裁剪:只选择需要的列,减少数据传输 -- 3. 连接重排序:优化连接顺序,先过滤再连接 -- 4. 局部聚合:在数据源层先进行部分聚合优化配置:
# 分布式查询优化配置 query_optimization: enabled: true # 谓词下推配置 predicate_pushdown: enabled: true supported_operators: ["=", ">", "<", ">=", "<=", "IN", "LIKE"] max_complexity: 10 # 最大谓词复杂度 # 列裁剪 column_pruning: enabled: true aggressive: false # 是否激进裁剪 # 连接优化 join_optimization: enabled: true algorithm: "cost_based" # 基于成本的优化 reorder_threshold: 8 # 最多重新排序8个表 # 聚合下推 aggregate_pushdown: enabled: true supported_functions: ["COUNT", "SUM", "AVG", "MIN", "MAX"] # 统计信息收集 statistics: auto_collect: true update_frequency: "1h" # 每小时更新一次 sample_rate: 0.1 # 10%的采样率6. AI功能实战:自然语言查询
6.1 WrenAI自然语言转SQL
WrenAI的核心功能是将自然语言转换为SQL查询,下面通过具体示例演示:
from wrenai import WrenAIClient import json # 初始化客户端 client = WrenAIClient( api_key="your_api_key", endpoint="https://api.wrenai.com/v1" ) # 自然语言查询示例 natural_language_query = """ 显示2024年第一季度每个月的销售总额,按月份排序 """ # 转换为SQL response = client.nl_to_sql( query=natural_language_query, data_source="sales_warehouse", context={ "tables": ["orders", "order_items", "products"], "business_glossary": { "销售总额": "SUM(oi.quantity * oi.unit_price)", "第一季度": "1月到3月" } } ) print("生成的SQL:") print(response.sql) print("\n解释:") print(response.explanation) print("\n置信度:", response.confidence) # 输出示例: """ 生成的SQL: SELECT DATE_TRUNC('month', o.order_date) as month, SUM(oi.quantity * oi.unit_price) as total_sales FROM orders o JOIN order_items oi ON o.order_id = oi.order_id WHERE o.order_date >= '2024-01-01' AND o.order_date < '2024-04-01' GROUP BY DATE_TRUNC('month', o.order_date) ORDER BY month; 解释: 这个查询从orders表获取订单日期,从order_items表计算每个订单项的销售额(数量×单价),然后按月份分组汇总,筛选2024年第一季度的数据。 置信度: 0.92 """6.2 复杂业务场景处理
对于复杂的业务需求,WrenAI能够理解业务逻辑并生成相应的SQL:
# 复杂业务查询示例 complex_query = """ 找出2024年购买金额超过10万元,但最近3个月没有下单的VIP客户, 显示客户姓名、最后下单日期和累计消费金额 """ response = client.nl_to_sql( query=complex_query, data_source="sales_system", context={ "business_rules": { "VIP客户": "累计消费金额 > 100000", "最近3个月": "最后下单日期 < CURRENT_DATE - INTERVAL '3 months'" }, "preferred_join_method": "INNER JOIN" } ) print("复杂查询SQL:") print(response.sql) # 生成的SQL可能类似: """ SELECT c.customer_name, MAX(o.order_date) as last_order_date, SUM(oi.quantity * oi.unit_price) as total_spent FROM customers c JOIN orders o ON c.customer_id = o.customer_id JOIN order_items oi ON o.order_id = oi.order_id WHERE o.order_date >= '2024-01-01' GROUP BY c.customer_id, c.customer_name HAVING SUM(oi.quantity * oi.unit_price) > 100000 AND MAX(o.order_date) < CURRENT_DATE - INTERVAL '3 months' ORDER BY total_spent DESC; """6.3 查询结果解释与可视化
WrenAI不仅生成SQL,还能解释查询结果并建议可视化方案:
# 执行查询并获取解释 query_result = client.execute_sql(response.sql) analysis = client.analyze_results(query_result) print("查询结果分析:") print(f"数据概览: 共{analysis.row_count}行,{analysis.column_count}列") print(f"关键洞察: {analysis.insights}") print("\n可视化建议:") for viz in analysis.visualization_suggestions: print(f"- {viz.chart_type}: {viz.reasoning}") if viz.example_config: print(f" 配置示例: {json.dumps(viz.example_config, indent=2)}") # 自动生成可视化配置 if analysis.visualization_suggestions: best_viz = analysis.visualization_suggestions[0] viz_config = client.generate_viz_config( data=query_result, chart_type=best_viz.chart_type, dimensions=best_viz.dimensions, measures=best_viz.measures ) print("\n生成的可视化配置:") print(json.dumps(viz_config, indent=2))7. 安全与权限管理
7.1 多层次安全架构
Canner/WrenAI提供企业级的安全保障,采用多层次安全架构:
# 安全配置示例 security: # 认证层配置 authentication: enabled: true providers: - type: "ldap" server: "ldap://company-ldap.example.com" base_dn: "dc=example,dc=com" - type: "oauth2" issuer: "https://auth.example.com" client_id: "canner-client" scopes: ["openid", "profile", "email"] - type: "saml" idp_metadata_url: "https://idp.example.com/metadata" sp_entity_id: "https://canner.example.com" # 授权层配置 authorization: model: "rbac" # 基于角色的访问控制 policies: - resource: "data_source:sales_mysql" actions: ["read", "query"] conditions: - "user.department == 'Sales'" - "time.between('09:00', '18:00')" - resource: "data_source:hr_postgres" actions: ["read"] conditions: - "user.role in ['HR', 'Manager']" # 数据保护层 data_protection: encryption: at_rest: true in_transit: true masking: enabled: true rules: - pattern: "*.phone" method: "partial" parameters: [3, 4, "****"] - pattern: "*.email" method: "hash" anonymization: enabled: true techniques: ["k-anonymity", "differential_privacy"]7.2 审计与合规性
满足企业审计和合规要求:
-- 审计日志表结构 CREATE TABLE audit_logs ( log_id BIGINT PRIMARY KEY, event_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP, user_id VARCHAR(100), user_ip INET, action VARCHAR(50), resource_type VARCHAR(50), resource_id VARCHAR(100), query_text TEXT, parameters JSONB, result_count INTEGER, execution_time_ms INTEGER, success BOOLEAN, error_message TEXT ); -- 创建审计策略 CREATE AUDIT POLICY data_access_audit ADD STATEMENTS (SELECT, INSERT, UPDATE, DELETE) ON DATABASE sales_db WHEN (1=1) -- 审计所有操作 LOG ALL; -- 查询审计日志 SELECT event_time, user_id, action, resource_type, resource_id, execution_time_ms, success FROM audit_logs WHERE event_time >= CURRENT_DATE - INTERVAL '7 days' AND action = 'SELECT' ORDER BY event_time DESC;审计报表生成:
import pandas as pd from datetime import datetime, timedelta class AuditReporter: def generate_compliance_report(self, start_date, end_date): """生成合规性报告""" query = """ SELECT user_id, COUNT(*) as total_queries, COUNT(CASE WHEN success = true THEN 1 END) as successful_queries, COUNT(CASE WHEN success = false THEN 1 END) as failed_queries, AVG(execution_time_ms) as avg_execution_time, MAX(execution_time_ms) as max_execution_time FROM audit_logs WHERE event_time BETWEEN %s AND %s GROUP BY user_id ORDER BY total_queries DESC """ df = pd.read_sql(query, self.connection, params=[start_date, end_date]) # 生成统计信息 report = { 'period': f"{start_date} 到 {end_date}", 'total_queries': df['total_queries'].sum(), 'unique_users': len(df), 'success_rate': (df['successful_queries'].sum() / df['total_queries'].sum() * 100), 'performance_metrics': { 'avg_execution_time': df['avg_execution_time'].mean(), 'max_execution_time': df['max_execution_time'].max() } } return report def generate_security_report(self): """生成安全审计报告""" # 检测异常访问模式 anomaly_query = """ WITH user_stats AS ( SELECT user_id, COUNT(*) as query_count, AVG(execution_time_ms) as avg_time, COUNT(DISTINCT resource_type) as unique_resources FROM audit_logs WHERE event_time >= CURRENT_DATE - INTERVAL '30 days' GROUP BY user_id ) SELECT user_id, query_count, avg_time, unique_resources, (query_count - avg_query_count) / stddev_query_count as z_score FROM user_stats CROSS JOIN ( SELECT AVG(query_count) as avg_query_count, STDDEV(query_count) as stddev_query_count FROM user_stats ) stats WHERE ABS((query_count - avg_query_count) / stddev_query_count) > 3 """ anomalies = pd.read_sql(anomaly_query, self.connection) return anomalies.to_dict('records')8. 运维监控与故障排查
8.1 系统监控配置
全面的监控体系确保平台稳定运行:
# Prometheus监控配置 scrape_configs: - job_name: 'canner' static_configs: - targets: ['canner-server:8080'] metrics_path: '/metrics' scrape_interval: 30s - job_name: 'canner_database' static_configs: - targets: ['postgres:5432'] metrics_path: '/metrics' # 关键监控指标 monitoring: metrics: # 系统资源指标 - name: "system_cpu_usage" query: "100 - (avg by (instance) (irate(node_cpu_seconds_total{mode='idle'}[5m])) * 100)" alert_threshold: 80 - name: "system_memory_usage" query: "(1 - (node_memory_MemAvailable_bytes / node_memory_MemTotal_bytes)) * 100" alert_threshold: 85 # 业务指标 - name: "active_queries" query: "canner_queries_active" alert_threshold: 100 - name: "query_duration_p95" query: "histogram_quantile(0.95, rate(canner_query_duration_seconds_bucket[5m]))" alert_threshold: 10.0 # 10秒 - name: "error_rate" query: "rate(canner_query_errors_total[5m]) / rate(canner_queries_total[5m]) * 100" alert_threshold: 5.0 # 5%错误率 # 告警规则 alerting: rules: - alert: "HighCPUUsage" expr: "system_cpu_usage > 80" for: "5m" labels: severity: "warning" annotations: summary: "CPU使用率过高" description: "实例 {{ $labels.instance }} 的CPU使用率持续5分钟超过80%" - alert: "QueryTimeout" expr: "query_duration_p95 > 30" for: "2m" labels: severity: "critical" annotations: summary: "查询响应时间过长" description: "P95查询响应时间超过30秒"8.2 日志管理与分析
集中式日志管理便于故障排查:
import logging import json from logging.handlers import RotatingFileHandler from elasticsearch import Elasticsearch class StructuredLogger: def __init__(self, app_name, log_level=logging.INFO): self.logger = logging.getLogger(app_name) self.logger.setLevel(log_level) # 结构化日志格式 formatter = logging.Formatter( '{"timestamp": "%(asctime)s", "level": "%(levelname)s", ' '"logger": "%(name)s", "message": "%(message)s", ' '"module": "%(module)s", "function": "%(funcName)s", ' '"line": "%(lineno)d"}' ) # 文件处理器 file_handler = RotatingFileHandler( f'/var/log/{app_name}.log', maxBytes=100*1024*1024, # 100MB backupCount=5 ) file_handler.setFormatter(formatter) self.logger.addHandler(file_handler) # Elasticsearch集成 self.es = Elasticsearch(['http://elasticsearch:9200']) def log_query(self, query_info): """记录查询日志""" log_entry = { 'timestamp': datetime.utcnow().isoformat(), 'level': 'INFO', 'type': 'query', 'query_id': query_info.get('query_id'), 'user_id': query_info.get('user_id'), 'data_source': query_info.get('data_source'), 'query_text': query_info.get('query_text'), 'execution_time': query_info.get('execution_time'), 'row_count': query_info.get('row_count'), 'status': query_info.get('status') } # 写入本地日志 self.logger.info(json.dumps(log_entry)) # 发送到Elasticsearch self.es.index(index='canner-query-logs', body=log_entry) def log_error(self, error_info): """记录错误日志""" error_entry = { 'timestamp': datetime.utcnow().isoformat