
Canner / WrenAI 全面解析新一代数据智能平台实战指南在日常数据开发工作中你是否遇到过这样的困境数据源分散在不同系统SQL编写效率低下业务人员难以自主分析数据传统的数据平台往往需要专业的数据工程师进行复杂配置业务团队的数据需求响应缓慢。本文将深入解析Canner和WrenAI这两个新一代数据智能平台从核心概念到实战应用带你掌握现代化数据平台的建设思路。1. 背景与核心概念1.1 什么是Canner和WrenAICanner和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之前需要确保环境满足基本要求。操作系统方面支持LinuxCentOS 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_TYPEpostgresql - CANNER_DB_HOSTpostgres - CANNER_DB_PORT5432 - CANNER_DB_NAMEcanner - CANNER_DB_USERcanner_user - CANNER_DB_PASSWORDyour_password depends_on: - postgres volumes: - ./config:/app/config postgres: image: postgres:13 environment: - POSTGRES_DBcanner - POSTGRES_USERcanner_user - POSTGRES_PASSWORDyour_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: fBearer {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, headersheaders, datajson.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_sourceds.name, statussuccess).set(duration) active_connections.labels(data_sourceds.name).set( get_active_connection_count(ds) ) except Exception as e: connection_errors.labels(data_sourceds.name).inc() query_duration.labels(data_sourceds.name, statuserror).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分钟缓存监控和管理APIimport requests import json class CacheManager: def __init__(self, base_url, auth_token): self.base_url base_url self.headers { Authorization: fBearer {auth_token}, Content-Type: application/json } def get_cache_stats(self): 获取缓存统计信息 response requests.get( f{self.base_url}/api/v1/cache/stats, headersself.headers ) return response.json() def clear_cache(self, cache_typeNone, patternNone): 清理缓存 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, headersself.headers, datajson.dumps(data) ) return response.json() def preload_cache(self, queries): 预加载缓存 data {queries: queries} response requests.post( f{self.base_url}/api/v1/cache/preload, headersself.headers, datajson.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自然语言转SQLWrenAI的核心功能是将自然语言转换为SQL查询下面通过具体示例演示from wrenai import WrenAIClient import json # 初始化客户端 client WrenAIClient( api_keyyour_api_key, endpointhttps://api.wrenai.com/v1 ) # 自然语言查询示例 natural_language_query 显示2024年第一季度每个月的销售总额按月份排序 # 转换为SQL response client.nl_to_sql( querynatural_language_query, data_sourcesales_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( querycomplex_query, data_sourcesales_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, indent2)}) # 自动生成可视化配置 if analysis.visualization_suggestions: best_viz analysis.visualization_suggestions[0] viz_config client.generate_viz_config( dataquery_result, chart_typebest_viz.chart_type, dimensionsbest_viz.dimensions, measuresbest_viz.measures ) print(\n生成的可视化配置:) print(json.dumps(viz_config, indent2))7. 安全与权限管理7.1 多层次安全架构Canner/WrenAI提供企业级的安全保障采用多层次安全架构# 安全配置示例 security: # 认证层配置 authentication: enabled: true providers: - type: ldap server: ldap://company-ldap.example.com base_dn: dcexample,dccom - 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 (11) -- 审计所有操作 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{modeidle}[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_levellogging.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, maxBytes100*1024*1024, # 100MB backupCount5 ) 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(indexcanner-query-logs, bodylog_entry) def log_error(self, error_info): 记录错误日志 error_entry { timestamp: datetime.utcnow().isoformat