// TABLE OF CONTENTS
  1. 数据管道整体架构设计
  2. 数据采集层设计与实现
  3. 数据存储层与ETL流程
  4. 分析层:指标计算与异常检测
  5. 可视化层:仪表盘与报告系统
  6. 管道运维与数据质量保障
CHAPTER 01

数据管道整体架构设计

GEO数据管道是将分散的GEO相关数据(AI搜索检测结果、网站分析数据、爬虫日志、内容发布记录等)统一采集、清洗、存储、分析并可视化的端到端系统。良好的数据管道是实现数据驱动GEO优化的基础设施。

五层架构设计:采集层(多种数据源的数据收集)、存储层(原始数据存储+处理后的分析数据存储)、处理层(数据清洗、转换、聚合)、分析层(指标计算、趋势分析、异常检测)、可视化层(仪表盘、报告、告警)。

技术选型原则:优先选择开源技术降低成本、选择团队熟悉的技术栈降低学习成本、选择有良好社区支持的技术降低维护风险、选择可水平扩展的技术应对数据增长。推荐技术栈:Python+Airflow(调度)+PostgreSQL(存储)+dbt(转换)+Metabase/Grafana(可视化)。

架构层 技术选型 核心职责 替代方案
采集层 Python脚本+API 多源数据收集 Logstash/Fluentd
存储层 PostgreSQL 结构化数据存储 MySQL/ClickHouse
处理层 dbt+Python 数据清洗转换 Spark/Airflow
分析层 Python+pandas 指标计算分析 SQL存储过程
可视化 Metabase+Grafana 仪表盘报告 Tableau/Looker
CHAPTER 02

数据采集层设计与实现

采集层负责从多个数据源收集GEO相关数据。核心数据源包括:AI搜索检测结果(Playwright/Selenium脚本采集)、网站分析数据(Google Analytics API/百度统计API)、爬虫日志(Nginx access log解析)、内容发布记录(CMS API/Webhook)、第三方监测数据(Profound API)。

采集层设计原则:每个数据源实现独立的采集器模块,统一输出为标准化的JSON格式;采集器支持定时调度和手动触发;采集失败时自动重试并记录错误日志;采集的数据附带元信息(数据源、采集时间、采集器版本)。

数据标准化:不同数据源的数据格式差异很大,采集层需要将数据统一为标准格式。标准数据结构包含:data_source(数据源标识)、entity_type(实体类型:keyword/page/crawler_log等)、timestamp(时间戳)、payload(具体数据内容)、metadata(元信息)。

example.py python
# 数据采集层:统一采集器框架
import abc
import json
from datetime import datetime

class BaseCollector(abc.ABC):
    '''数据采集器基类'''
    def __init__(self, source_name):
        self.source_name = source_name

    @abc.abstractmethod
    def collect(self, **kwargs):
        '''执行数据采集,返回标准格式数据列表'''
        pass

    def _format_record(self, entity_type, payload, metadata=None):
        '''格式化为标准数据结构'''
        return {
            'data_source': self.source_name,
            'entity_type': entity_type,
            'timestamp': datetime.now().isoformat(),
            'payload': payload,
            'metadata': metadata or {}
        }

class GoogleAnalyticsCollector(BaseCollector):
    '''Google Analytics数据采集器'''
    def __init__(self, ga_credentials):
        super().__init__('google_analytics')
        self.credentials = ga_credentials

    def collect(self, start_date, end_date, metrics=['sessions', 'pageviews']):
        '''从GA API采集流量数据'''
        # 模拟GA API调用
        # 实际使用时替换为google-analytics-data API调用
        data = []
        for date in self._date_range(start_date, end_date):
            record = self._format_record(
                entity_type='traffic_daily',
                payload={
                    'date': date,
                    'sessions': 1200,  # 实际从API获取
                    'pageviews': 3500,
                    'ai_search_visits': 85  # AI搜索来源流量
                }
            )
            data.append(record)
        return data

    def _date_range(self, start, end):
        from datetime import datetime, timedelta
        d1 = datetime.strptime(start, '%Y-%m-%d')
        d2 = datetime.strptime(end, '%Y-%m-%d')
        dates = []
        while d1 <= d2:
            dates.append(d1.strftime('%Y-%m-%d'))
            d1 += timedelta(days=1)
        return dates

class CrawlerLogCollector(BaseCollector):
    '''爬虫日志采集器'''
    def __init__(self, log_path):
        super().__init__('crawler_log')
        self.log_path = log_path

    def collect(self, date=None):
        '''解析Nginx日志,提取AI爬虫活动'''
        data = []
        log_file = f'{self.log_path}/access_{date or datetime.now().strftime("%Y%m%d")}.log'
        try:
            with open(log_file, 'r', encoding='utf-8') as f:
                for line in f:
                    # 使用第75章的日志解析逻辑
                    parsed = parse_nginx_log(line)
                    if parsed and parsed.get('crawler'):
                        record = self._format_record(
                            entity_type='crawler_request',
                            payload=parsed
                        )
                        data.append(record)
        except FileNotFoundError:
            print(f'日志文件不存在: {log_file}')
        return data
CHAPTER 03

数据存储层与ETL流程

存储层负责持久化采集到的数据。推荐使用PostgreSQL作为主数据库,因为它对JSON数据类型有原生支持,且具有强大的分析查询能力。对于大规模时序数据(如爬虫日志),可以额外使用ClickHouse作为分析仓库。

数据库表设计:原始数据表(raw_data,存储采集层输出的标准格式数据)、关键词维度表(dim_keywords,存储关键词及其分类和优先级)、日期维度表(dim_dates,存储日期及其属性)、事实表(fact_daily_metrics,存储每日聚合的GEO效果指标)。

ETL(Extract-Transform-Load)流程:Extract阶段从原始数据表中提取当日数据;Transform阶段进行数据清洗(去重、缺失值处理)、聚合计算(按关键词×平台×日期聚合计算提及率等指标)、维度关联(关联关键词分类信息);Load阶段将处理后的数据写入事实表。

ETL调度:使用Airflow或简单的cron调度每日ETL任务。ETL任务应在数据采集完成后执行(建议凌晨2点运行,确保当日数据已全部采集完毕)。ETL执行失败时自动重试并发送告警。

example.py python
# ETL流程:数据清洗与聚合
import psycopg2
from psycopg2.extras import execute_batch
import pandas as pd
from datetime import datetime, timedelta

class GEODataPipeline:
    def __init__(self, db_config):
        self.conn = psycopg2.connect(**db_config)

    def run_daily_etl(self, target_date=None):
        '''执行每日ETL流程'''
        date = target_date or (datetime.now() - timedelta(days=1)).strftime('%Y-%m-%d')
        print(f'开始ETL流程: {date}')

        # Step 1: Extract - 提取当日原始数据
        raw_data = self._extract(date)
        print(f'  提取数据: {len(raw_data)} 条')

        # Step 2: Transform - 数据清洗与聚合
        metrics = self._transform(raw_data, date)
        print(f'  聚合指标: {len(metrics)} 条')

        # Step 3: Load - 写入事实表
        self._load(metrics)
        print(f'  写入完成: {date} ETL完成')

    def _extract(self, date):
        '''提取当日原始数据'''
        query = '''
        SELECT payload->>'keyword' as keyword,
               payload->>'platform' as platform,
               (payload->>'brand_mentioned')::boolean as brand_mentioned,
               payload->>'response' as response,
               timestamp
        FROM raw_data
        WHERE entity_type = 'ai_search_result'
          AND DATE(timestamp) = %s
        '''
        df = pd.read_sql(query, self.conn, params=(date,))
        return df

    def _transform(self, df, date):
        '''数据清洗与聚合计算'''
        if df.empty:
            return []

        # 按关键词×平台聚合
        agg = df.groupby(['keyword', 'platform']).agg(
            total_checks=('brand_mentioned', 'count'),
            brand_mentions=('brand_mentioned', 'sum')
        ).reset_index()
        agg['mention_rate'] = agg['brand_mentions'] / agg['total_checks']
        agg['date'] = date

        return agg.to_dict('records')

    def _load(self, metrics):
        '''写入事实表'''
        query = '''
        INSERT INTO fact_daily_metrics
            (date, keyword, platform, total_checks, brand_mentions, mention_rate)
        VALUES (%(date)s, %(keyword)s, %(platform)s, %(total_checks)s, %(brand_mentions)s, %(mention_rate)s)
        ON CONFLICT (date, keyword, platform) DO UPDATE
        SET total_checks = EXCLUDED.total_checks,
            brand_mentions = EXCLUDED.brand_mentions,
            mention_rate = EXCLUDED.mention_rate
        '''
        with self.conn.cursor() as cur:
            execute_batch(cur, query, metrics)
        self.conn.commit()
CHAPTER 04

分析层:指标计算与异常检测

分析层基于事实表中的聚合数据,计算GEO核心指标并进行异常检测。分析层的输出直接驱动可视化层和告警系统。

核心指标计算:整体品牌提及率(所有关键词×平台的加权平均提及率)、分平台提及率(各AI平台的平均提及率)、分关键词类别提及率(品牌词/行业词/长尾词的分别提及率)、提及率趋势(7日/30日移动平均)、竞品对比指标(品牌vs竞品的提及率差异)。

异常检测算法:使用Z-Score方法检测提及率异常波动。计算过去30天提及率的均值和标准差,当当日提及率偏离均值超过2个标准差时标记为异常。对于新发布的内容页面,使用首次被AI引用的时间间隔作为新鲜度指标,超过72小时未被引用则触发Info级别提醒。

分析任务调度:核心指标每日计算一次(ETL完成后自动触发),趋势分析每周计算一次,异常检测每日运行并在检测到异常时实时告警。所有分析结果写入分析结果表供可视化层查询。

example.py python
# 分析层:指标计算与异常检测
import numpy as np
from scipy import stats

class GEOAnalyzer:
    def __init__(self, db_config):
        self.db_config = db_config

    def calculate_core_metrics(self, date):
        '''计算核心GEO效果指标'''
        conn = psycopg2.connect(**self.db_config)

        # 整体提及率
        query = 'SELECT AVG(mention_rate) as overall_rate FROM fact_daily_metrics WHERE date = %s'
        overall = pd.read_sql(query, conn, params=(date,)).iloc[0]['overall_rate']

        # 分平台提及率
        query = 'SELECT platform, AVG(mention_rate) as rate FROM fact_daily_metrics WHERE date = %s GROUP BY platform'
        by_platform = pd.read_sql(query, conn, params=(date,))

        # 7日移动平均
        query = '''SELECT date, AVG(mention_rate) as daily_rate
                   FROM fact_daily_metrics
                   WHERE date >= %s::date - 7
                   GROUP BY date ORDER BY date'''
        trend = pd.read_sql(query, conn, params=(date,))
        trend['ma7'] = trend['daily_rate'].rolling(7, min_periods=1).mean()

        conn.close()

        return {
            'date': date,
            'overall_mention_rate': float(overall) if overall else 0,
            'by_platform': by_platform.to_dict('records'),
            'trend': trend.to_dict('records')
        }

    def detect_anomalies(self, date, lookback_days=30):
        '''异常检测:Z-Score方法'''
        conn = psycopg2.connect(**self.db_config)
        query = '''
        SELECT date, AVG(mention_rate) as daily_rate
        FROM fact_daily_metrics
        WHERE date >= %s::date - %s AND date < %s
        GROUP BY date ORDER BY date
        '''
        historical = pd.read_sql(query, conn, params=(date, lookback_days, date))

        # 获取当日数据
        query_today = 'SELECT AVG(mention_rate) as rate FROM fact_daily_metrics WHERE date = %s'
        today_rate = pd.read_sql(query_today, conn, params=(date,)).iloc[0]['rate']
        conn.close()

        if today_rate is None or len(historical) < 7:
            return {'anomaly': False, 'reason': '数据不足'}

        mean = historical['daily_rate'].mean()
        std = historical['daily_rate'].std()
        z_score = (today_rate - mean) / std if std > 0 else 0

        return {
            'anomaly': abs(z_score) > 2,
            'z_score': float(z_score),
            'today_rate': float(today_rate),
            'historical_mean': float(mean),
            'historical_std': float(std),
            'direction': '下降' if z_score < -2 else '上升' if z_score > 2 else '正常'
        }
CHAPTER 05

可视化层:仪表盘与报告系统

可视化层将分析结果转化为直观的图表和报告,帮助GEO团队和管理层快速理解数据含义并做出决策。推荐使用Metabase(开源、易用、支持SQL查询)作为主仪表盘工具,Grafana用于实时监控指标。

Metabase仪表盘设计:核心指标卡片区(整体提及率、AI搜索流量、AI引用页面数、竞品对比胜率)、趋势折线图区(7日/30日提及率趋势、分平台趋势对比)、关键词分析区(提及率最高/最低的关键词排行、关键词覆盖热力图)、平台对比区(各AI平台提及率柱状图、平台×关键词矩阵热力图)。

自动报告生成:使用Python脚本每周自动生成GEO效果周报,包含核心指标摘要、趋势分析、异常事件记录、优化建议。报告以PDF格式通过邮件发送给管理层和GEO团队。

数据API层:为支持自定义分析和第三方集成,搭建RESTful API层提供数据查询接口。API支持按日期范围、平台、关键词等维度过滤查询,返回JSON格式数据。API使用FastAPI框架实现,支持API Key认证和速率限制。

CHAPTER 06

管道运维与数据质量保障

数据管道的稳定运行和 数据质量保障是长期工程。需要建立完善的监控、告警和质量检查机制,确保管道持续产出可靠的数据。

管道监控:监控每个ETL任务的执行状态(成功/失败/运行中)、执行时长、处理数据量。使用Airflow的DAG监控面板或自建监控页面展示管道运行状态。任务失败时自动重试(最多3次),重试仍失败则发送Critical级别告警。

数据质量检查:每个ETL任务完成后自动执行数据质量检查,包括完整性检查(关键字段是否为空)、一致性检查(数据总量是否在合理范围)、准确性检查(提及率是否在0-100%之间)、时效性检查(数据是否为当日数据)。质量检查不通过时阻止数据写入事实表并发送告警。

数据血缘追踪:记录数据从采集到最终展示的完整流转路径,便于问题排查。每条数据记录附带数据血缘信息(来源、处理步骤、处理时间)。当发现数据异常时,可以通过数据血缘快速定位问题环节。

运维最佳实践

ETL任务在凌晨2-4点执行,避开业务高峰期

数据库每周全量备份,每日增量备份

保留原始数据至少6个月,聚合数据至少2年

每月进行一次数据质量审计,检查数据完整性和准确性

每季度评估管道性能,优化慢查询和瓶颈环节