ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

AWS上自建CDP全指南:架构选型、身份合并与避坑实践

AWS上自建CDP全指南:架构选型、身份合并与避坑实践 简介《智能客户数据平台的AWS云端之旅》是一份基于亚马逊云服务构建客户数据平台CDP的解决方案演示文稿面向企业数据团队、营销运营人员及方案架构师。内容从CDP定义切入梳理企业7×24小时稳定运行、弹性扩容、降低运维成本等需求给出以EC2、EMR、S3、CloudFront为核心的技术选型并拆解实时与非实时数据采集、清洗、ID打通、标准化到360度画像的完整链路。PPT虽仅1个pptx文件、大小1.66MB却覆盖AI驱动的全客户生命周期管理如新客获取、客户提升、成熟、衰退各阶段以及RFM模型、流失预警、Look-alike、营销计分等关键算法同时包含官网、APP、微信、门店等多渠道应用场景以及某信用卡中心个性化推荐、微信运营优化的落地案例。对希望理解CDP与AWS结合方式、快速搭建方案框架的读者来说这份精简的材料可作为评审汇报或技术交流的参考。目前已有152人学习浏览。1. 智能客户数据平台迁到 AWS为什么说这事能成、坑在哪做营销或用户增长的同学大概率见过这种场景用户行为在埋点系统里订单在电商库里客服记录在 CRM 里每次活动前都要求数据团队提数、清洗、拉宽表做完活动发现有漏数据又要重来一遍。客户数据平台CDPCustomer Data Platform就是把散落在各处的客户数据汇到一个地方统一身份、算标签、做分群让运营能自己查、自己圈人。而「AWS 云端之旅」讲的是这套 CDP 不买商业套件、不自己租机房而是用 AWS 上的一堆托管服务把它搭出来。这条路能省运维人力但它不是把数据丢进 S3 就完事中间有大量设计决策比如身份合并怎么做、明细和画像怎么分层、成本怎么控。这篇文章按我实际搭过的路径讲一遍从选型到建管道到分群上线连同几个能让你少加班的参数设置。2. 方案选型与架构AWS 上搭 CDP 的六层骨架和三个选型理由2.1 不买 Salesforce CDP、不折腾开源套件选 AWS 自建的判断逻辑先聊一个容易纠结的问题市面上有现成的 CDP 产品Salesforce CDP、Segment、Amplitude也有开源的 ClickHouse Kafka 自建方案为什么还有人选 AWS 自建常见理由无非是成本、数据合规、二次开发空间但放到实际项目里最硬的理由往往是一个数据已经有一部分在 AWS 上了。比如埋点日志本来就打到 S3订单库在 RDS恰好又上了 Redshift 做 BI那再加一层 CDP 时与其再引入一套商业 CDP 把数据拉出去走一遍管道不如在现有 AWS 体系里搭一层逻辑模型让管道更短、权限边界更清晰。另一个常被忽略的点是数据隐私协议的响应时间。商业 CDP 一般会要求你把原始事件往外送一份这会在合规评审时多一道口子自建 CDP 的好处是原始数据可以留在原有存储区域只向上层开放加工后的画像和标签。AWS 的架构让这种「数据不出域」可以通过账号隔离或 S3 Bucket Policy 直接实现不需要跟厂商反复确认对方的数据处理协议。当然自建也意味着数据管道、任务调度、质量监控都得自己负责所以这个选择更适合有 1 到 2 名数据工程师、且不想被商业套件的 Schema 绑死的团队。2.2 用哪些 AWS 服务搭 CDP组件清单与职责边界我搭过的 AWS CDP 架构大致分六层这里先给你一张组件和职责对照表后文每一步都会用到它。层级使用的 AWS 服务在本方案中的职责接入层Kinesis Firehose, AppFlow, Lambda收取埋点事件、拉取 SaaS 数据、轻量清洗存储层S3原始区/明细区/画像区原始事件、明细宽表、画像快照的分区存储加工层AWS Glue ETL, Glue Studio做 ID 合并、标签计算、分群聚合服务层RDS 或 Redshift存放明细数据与画像结果供实时查询分析层QuickSight画像分析、标签覆盖率看板编排层EventBridge Step Functions任务调度、依赖触发与失败告警这套选型和直接买 ClickHouse 自建相比最直接的好处是不用养一个 ClickHouse 集群。Glue 和 Firehose 都是按量付费数据量涨了不用提前扩容夜里跑 ETL 贵一点但可接受。Redshift 可以换成 RDS for PostgreSQL 起步数据量在千万行级别、查询并发不高时用 RDS 更省钱而 QuickSight 连 RDS 做标签分析完全够用。不过这里有一个坑需要提前说AWS 组件选型容易只按「数据要流进来」的直觉排忽略 CDP 最核心的身份合并需要反复刷历史数据。如果只用 Firehose Glue 做一次性接入后续要修改身份映射规则比如把「手机号优先」改成「邮箱优先」时历史数据不会自动重算必须在 Step Functions 里加一个重跑分区的前置任务否则标签结果会对不上。2.3 初始数据模型明细、事件、画像三类表的分层逻辑CDP 的数据模型和普通数仓宽表不太一样。普通宽表把一行用户的所有属性放一起CDP 则必须把原始事实和加工结果分开因为一个用户可能有多个设备、多个账号、多个邮箱身份合并这件事随时会变。我的做法是建三类核心表对应 S3 的三个 prefix第一类是raw_events存原始行为事件schema 几乎是埋点上报的原文包括 event_id、user_id可能为空、anonymous_id、event_type、event_time、propertiesmap 类型不加工、不合并保留全部字段。第二类是customer_profile这是画像主表每行是一个统一后的客户包含 unified_id、手机号、邮箱、默认地址、生命周期阶段、最近一次活跃时间等。第三类是customer_tags一行一个标签键值对例如tag_name高活跃tag_valuetrue这种设计比在 profile 表里加一堆布尔列灵活运营圈人选标签时只需要一条 SQL 过滤tags表。注意customer_profile和customer_tags的主键都是unified_id但unified_id不是埋点原生的。它是在加工过程中由 ID 映射规则生成的所以在这两张表里看到 ID 到了千万级之后一定要给unified_id建排序键Redshift 的 SORTKEY或索引RDS。有次我在这张表上用WHERE unified_id ...查一个用户没加索引前跑了 8 秒加上之后秒回——这类查询在运营接的后台里频率很高不加索引会是慢性病。3. 数据采集与存储设计把散落的用户数据收进 S33.1 搭一条最小可用的事件采集管道Kinesis Firehose S3不管你的数据源是 Web 埋点、App SDK 还是服务端日志AWS 上最省事的接法是 Kinesis Firehose 直写 S3。先创建一个 Firehoseaws firehose create-delivery-stream \ --delivery-stream-name cdp-user-event-stream \ --s3-destination-configuration \ BucketARNarn:aws:s3:::your-cdp-bucket \ RoleARNarn:aws:iam::123456789012:role/firehose-s3-role \ Prefixraw/events/year!{timestamp:yyyy}/month!{timestamp:MM}/day!{timestamp:dd}/ \ BufferingHints{SizeInMBs64,IntervalInSeconds300} \ CompressionFormatGZIP这段命令做了一件事把实时事件流落到 S3 的raw/events目录下按天分区gzip 压缩。两个 BufferingHints 参数值得说清楚SizeInMBs64表示等攒到 64 MB 才落一次盘IntervalInSeconds300表示最多等 300 秒两者先到先触发。新手会纠结要不要调小一点好让数据更实时我的建议是不要小于 32 MB / 300 秒因为 S3 对单个前缀的写入性能有限制小文件一多后面 Glue 跑批会被文件列表压垮。采集层必须注意数据完整性。Firehose 写到 S3 是至少一次语义理论上会有重复事件。处理手段是在下游 Glue 里按event_id event_time做去重而不是指望采集层不重。另外raw_events表最好按year/month/day三层分区在建表时就定好不要用date单分区因为后续数据回溯和补数会按天刷多层分区能让你刷某一天时不触碰整月数据。3.2 建明细层表结构一次到位避免返工事件落到 S3 后需要用 Glue 把它解析成列式存储供后续 ETL 和查询使用。这里建议直接建一个 Iceberg 或 Hudi 表而不是静态 Parquet 分区表。原因很简单CDP 的数据会反复更新你需要支持数据回刷。我第一次做时用了parquet partition的静态表后来改 ID 映射规则历史分区要重写还得手动删分区非常痛苦。用 Iceberg 表后Glue 的merge into能直接覆盖指定分区的数据。CREATE TABLE cdp.raw_events_iceberg ( event_id STRING, user_id STRING, anonymous_id STRING, event_type STRING, event_time TIMESTAMP, properties MAPSTRING, STRING ) PARTITIONED BY (dt STRING) LOCATION s3://your-cdp-bucket/iceberg/raw_events/ TBLPROPERTIES (table_typeICEBERG)这条 SQL 是在 Glue Data Catalog 里建 Iceberg 表的典型写法。PARTITIONED BY (dt STRING)里的dt是事件日期取自上一节 Firehose 目录里的year/month/day在 Glue 任务里通过一个date_format(event_time, yyyy-mm-dd)生成。properties用MAP而不是把埋点字段全部单独设列是为了兼容未来埋点迭代——埋点新增字段只需要加 MAP 的 key不需要改表结构。建表之后拿什么验证你可以先跑一条查询看看最近一天的事件数和独立 anonymous_id 数SELECT dt, count(*) AS event_cnt, approx_distinct(anonymous_id) AS uid_cnt FROM cdp.raw_events_iceberg WHERE dt current_date - interval 7 day GROUP BY dt;approx_distinct是一个近似去重函数在千万级事件上比count(distinct)快很多。看到每条设备的匿名用户数稳定、不出现断崖下跌就说明采集管道是通的。这里有个常见坑Firehose 会把失败的记录写进同一个 S3 bucket 的processing-failed/前缀下很多人以为管道没报错就没问题结果某天查数发现少了一段这是第一道要盯的检查点。3.3 批式数据源接入AppFlow 与 JDBC 连接 RDS埋点事件是流式的但订单、CRM、客服记录一般都在业务数据库里需要批量同步。AWS 的 AppFlow 支持从 Salesforce、SAP、Zendesk 等 SaaS 拉数据到 S3适合没有技术团队维护接口的源如果数据在 RDS 或自建 MySQL 里我一般直接写 Glue 作业用 JDBC 读取。Glue 从一个 RDS PostgreSQL 的orders表同步订单最小代码是这样from awsglue.context import GlueContext from awsglue.job import Job from awsglue.utils import resolveChoice from pyspark.context import SparkContext sc SparkContext() glueContext GlueContext(sc) spark glueContext.spark_session source_df glueContext.create_dynamic_frame.from_catalog( database cdp, table_name rds_orders ) cleaned_df source_df.drop_fields([internal_note, raw_response]) cleaned_df cleaned_df.rename_field(cust_phone, phone) cleaned_df cleaned_df.rename_field(order_total_cents, order_amount_cents) glueContext.write_dynamic_frame.from_catalog( frame cleaned_df, database cdp, table_name cdp_orders, transformation_ctx datasink ) job.commit()这段代码看起来简单但三个细节值得说明。第一rename_field(order_total_cents, order_amount_cents)是把源库里的字段名对齐到 CDP 明细层的统一命名CDP 的字段命名必须对所有数据源统一否则后面做标签时会为「到底用 total 还是 amount」吵个不停第二drop_fields([internal_note, raw_response])在做列裁剪同步明细表时没必要把业务表里的内部备注也拉进来省 S3 存储也降低 PII 泄露风险第三这里没有写filter是因为全量同步更容易让下游判断数据是否完整做增量可以依靠 Glue 的 bookmarks 功能但初次搭建不建议同时开增量先全量跑通再优化调度。这一步最值得做的验证是把同步前后的行数对比接到告警里Glue 任务结束前用job.commit()前的计数或者 CloudWatch 里看glue.driver.ExecutorAllocationManager相关指标。我见过不止一次源业务库某张表被误删了分区同步任务不报错、只是行数变少导致标签数据悄悄失真。行数波动告警是 CDP 管道里的第一道安全防线。4. 用户身份合并与画像构建从 IDs 到 unified_id 的关键一跳4.1 ID-Mapping为什么用户表不能只用手机号做主键CDP 的核心不是「有多少数据」而是「能不能把同一个人在不同设备、不同账号下的行为串起来」。一个人可能在小程序里用微信登录、在 App 里用手机号注册、在 PC 网页里没登录只留下了 cookie 设备 ID这三条记录单独看是三个人实际上是一个客户。把这三条记录归到同一个unified_id的过程就是 ID-Mapping也是 CDP 建设里最值得慢下来的地方。我的做法是用规则引擎分步合并而不是靠机器学习。第一步把存在确定性标识手机号或邮箱的日志归并比如identity_phone、identity_email第二步再把匿名设备 ID 通过「同设备登录过账号」这样的行为关联挂到第一步的 unified_id 下。场景上一个anonymous_id如果曾经在同一设备上绑定过某个已登录用户就把它归属进去。实现上用一个全局映射表id_mappingCREATE TABLE cdp.id_mapping ( unified_id STRING, id_type STRING, id_value STRING, first_seen_at TIMESTAMP, last_seen_at TIMESTAMP )这个表里一条记录代表「某个 id_value 归属于某个 unified_id」id_type可以是phone、email、anonymous_id、user_id等。跑 ID-Mapping 时先把所有 id_type 按规则两两 join合并成候选对再找连通分量。连通分量的计算是个计算密集活用 Spark 的graphx或 Python 的networkx都可以做但需要注意数据规模会检验你的实现。如果团队里没有强图计算经验我建议先用一个简单方案起步分层合并不追求全图连通。也就是先按手机号合并再按月跑一轮「手机号 邮箱」再合并优先级是可解释性大于每轮合并率。每轮合并后的结果写回id_mapping表并记录合并规则版本号。比起一步到位的图算法这种做法可回滚、可解释运营同事问「这两个 ID 怎么归到一起的」时你能拿出规则去解释而不是丢给一个黑匣子。4.2 用 Glue 实现增量合并回溯一周的实体合并ID-Mapping 的批量任务我一般放在日批里每天回溯最近 7 天的数据避免长时间不回看导致同人分裂成两个 unified_id。核心 ETL 逻辑是from pyspark.sql import functions as F # 读取最近7天的原始事件找登录事件里带user_id的记录 login_df spark.sql( SELECT anonymous_id, user_id, event_time FROM cdp.raw_events_iceberg WHERE dt date_format(current_date - interval 7 day, yyyy-MM-dd) AND event_type user_login AND user_id IS NOT NULL ) # 统计每个anonymous_id最近7天里登录过的user_id去重列表 mapping_df login_df.groupBy(anonymous_id) \ .agg(F.collect_set(user_id).alias(user_id_list), F.max(event_time).alias(last_login_time)) # 如果anonymous_id只对应一个user_id则绑定到该user_id bind_df mapping_df.filter(F.size(user_id_list) 1) \ .withColumn(id_type, F.lit(anonymous_id)) \ .withColumn(user_id, F.element_at(F.col(user_id_list), 1)) \ .select(user_id, id_type, anonymous_id, last_login_time)这个片段的逻辑是只用登录事件做确定性绑定。collect_set聚合出最近 7 天内每个设备登录过的用户 ID 列表然后只绑定那些唯一对应的记录如果一个设备登录过多个账号家庭共用设备很常见就不绑定留到下一轮规则处理。我不建议一上来就把共用设备硬绑定到最近登录的账号那会把夫妻俩的数据混到一个画像里后患很大。执行完之后把bind_df写入id_mapping同时要更新订单表和事件表里的 user_id 字段使它们指向 unified_id。这一步常见做法是对订单表做JOIN id_mapping ON orders.user_id id_mapping.user_id再UPDATE orders SET unified_id id_mapping.unified_id。注意这个 UPDATE 不要在全量数据上做只对orders.dt 最近 7 天的分区做否则全表重写会非常昂贵。4.3 聚合出客户画像把行为归纳成标签并用 SQL 验证覆盖率身份合并跑完画像表就有了主键基础。接下来把行为数据聚合成 profile 和 tags。这个环节最容易翻车的是「先定标签后建管道」比如运营要一个「高活跃用户」标签但「活跃」的定义是每周登录一次还是每周下单一次必须提前书面定好。否则千辛万苦算完之后运营说不是这个口径重跑成本极高。标签计算的 SQL 示例以「近 30 天下单次数」为例INSERT INTO cdp.customer_tags (unified_id, tag_name, tag_value, updated_at) SELECT o.unified_id, order_cnt_30d AS tag_name, cast(count(distinct o.order_id) as string) AS tag_value, now() AS updated_at FROM cdp.customer_orders o WHERE o.order_time now() - interval 30 day GROUP BY o.unified_id这个 SQL 在做的是把订单明细聚合到 unified_id 上然后写入标签宽表。每一个标签就是 INSERT 一行tag_name和tag_value分开方便后续筛选和做标签管理页面。计算标签是简单活但要注意count(distinct o.order_id)在数据量大时会很慢尤其是日订单量过百万时。可以先对订单明细按(unified_id, order_id)去重生成子表再 count而不是两个 distinct 直接怼上去。画像表验证可以做三个维度第一是unified_id总数与源系统里已知客户数的覆盖差是否在预期范围第二是重点账号抽查拿一个已知测试手机号看看他在customer_tags里的标签是否和他在源库的实际行为一致第三是空值率检查比如「近 30 天首次登录时间」这个标签至少 80% 的用户不该为空。这三个检查要固化成 SQL塞进 Step Functions 里作为每天 ETL 结束后的一次质量快照通过 CloudWatch 报警。5. CDP 云上落地避坑指南五个高频故障和对应的修复路径5.1 小文件过多把 Glue 批处理拖垮现象是数据量并不大但 Glue 作业跑了两个多小时日志里大量 time 在 S3 文件列表操作和 task 调度上。原因是 Firehose 的 Buffer 设置过小默认 5 MB 或者有人为了追求「实时」把它调成 1 MB导致 S3 里堆了几百万个小文件Glue 做一次简单的count都要先遍历一遍小文件列表。解决路径是三层同时调整Firehose 的BufferingHints调大到 64 MB / 300 秒Glue 作业开启groupFiles和groupSize在 Spark 里设置spark.files.openCostInBytes和spark.files.maxPartitionBytes让 Spark 一次多读几个文件同时给 Glue 作业加一个定期OPTIMIZE的 compaction 任务把 Iceberg 表的文件数降下来。这个做完之后同样的数据量从两个小时降到了十几分钟。5.2 ID 映射规则调整导致历史标签全盘作废现象是运营发现某个分群的用户数突然掉了一半回查发现是上次改造了匿名设备绑定规则。原因是 ID-Mapping 是每天跑批的但标签表是 T1 生成的规则变了之后只有新日期分区的标签用了新规则历史分区的标签还是旧的给人感觉就是数据错了。也就是说ID-Mapping 规则是全局有状态的逻辑改变规则必须重刷历史分区。解决方法是先为每一批次的id_mapping和customer_tags写入带上一个rule_version字段下次改规则时明确重刷涉及的分区比如「把近 90 天的标签全部重算」再用MERGE INTO覆盖更新这些分区而不是 INSERT。Step Functions 里要设计一个「重跑指定分区」的手动触发任务。吃一堑之后我做的第一个改造就是加 rule_version这比几个月后被人追着问「为什么这个用户上月还存在、这月消失了」要省心得多。5.3 Redshift 查询并发一高就排队运营后台报表大量超时现象是 QuickSight 或者内部运营后台一刷标签列表Redshift 出现大量WLM queue wait查询一篇一片超时。原因是查询 SQL 直接打到了 Redshift 的生产队列上和 ETL 抢资源。CDP 的标签查询模式是典型的高频小查询而不是大分析所以不应该跟跑批抢同一个队列。解决方法是把 Redshift 的 WLM 配置拆成三部分一个etl_queue固定 50% 内存给 Glue 批作业用一个dashboard_queue给 QuickSight限制并发数并设置max_execution_time一个ad-hoc_queue给分析师临时查询。同时在 QuickSight 里开 SPICE 加速把标签表缓存到内存里而不是每次查询都打 Redshift。这样设置之后运营的页面基本能做到 1~2 秒出数而不是每次点按钮后去喝咖啡。5.4 S3 权限策略覆盖不全出现「部分人能查、部分人看不到」的诡异问题现象是这个角色能查到某个标签另外一个角色死活查不到但看了一遍 Lake Formation 权限又是好的。原因是 Glue 作业写数据时的 location 权限和查询时用的表权限是两套体系。如果之前给的是「某些 prefix 可写但没读」下游角色读表时就会遇见元数据能看到、数据拿不到的错位局面。解决方法是统一用Lake Formation做数据湖权限把raw / processed / profile / tags四个库分别授权按「读」「写」「读写」区分角色不要再用单独的 S3 Policy 逐个前缀去配。另外每一次建新表之后要复查一次 location 权限因为新表会自动生成一个湖里新路径很容易漏配。5.5 埋点属性新增字段后历史数据解析失败现象是某个新版本的埋点把properties里的某个字段类型从字符串改成了对象Glue 解析任务在某个历史分区报错导致全部下游不更新。原因是 Iceberg 表的 schema 是强类型新增字段的类型和旧数据不一致。解决方法是不要直接改原字段类型而是在 Glue 里做resolveChoice转换统一 cast 成 string或者把新字段单独存为一个新 key比如properties.amount_v2在 ETL 里显式转换。这里要记住一个原则CDP 的原始层应该是「存储原样」加工层负责兼容不要让上游埋点来迁就表的 schema。6. 进阶成本治理与画像质量验证的四个日常动作走到这一步 CDP 已经能跑起来了接下来值得投入的是成本和质量的长期治理。先看成本AWS CDP 的主要开销通常集中在 S3 存储、Glue DPU 小时数和 Redshift 或 RDS 的实例费。S3 可以用生命周期规则把 180 天前的原始事件自动转GLACIER存储冷数据不影响查询但单价大幅下降筛选条件是原始事件表基本不会用来做 T30 以上的分析更多是合规存档和问题追溯才需要。转换不要设得太激进建议原始区保留 180 天标准存储画像区保留 90 天标准存储再老的全部转冷。Glue 成本的优化一般从作业类型入手。每天跑批的标签任务如果对延迟不敏感可以改用G.1X而非G.2X的 DPU 类型并且把几个小时级任务合并到一个 Glue 作业里串行跑而不是拆成十个作业分别起十个集群。我见过最夸张的浪费是一个团队的 21 个 Glue 作业每小时各起一个 10 DPU 的集群合并后成本降了接近一半。另外给 Glue 作业设置合理的超时时间防止死循环跑一夜扣费。画像质量验证方面我建议每周跑一次一致性检查。具体做法是准备一份「金标准」测试账号表里面包含十几个自定义的测试用户他们行为已知ETL 跑完后对这些测试账号逐一比对标签值是否与预期一致任何一条不匹配都触发告警。这套机制比检查整个表的统计值更早暴露规则问题。第二个检查是统一 ID 覆盖率即id_mapping表里每天有事件的 anonymous_id 中有多少比例最终被绑定到一个 unified_id。这个比例如果连续三天下降说明 ID-Mapping 的新规则有漏洞或者新来的流量渠道没有被覆盖。第三个检查是重复率抽查 1000 条 user_id 看是否映射到同一个 unified_id这个数字超过 5% 就要查是不是测试环境或者多租户数据跑进生产了。最后一个建议是把 CDP 的运营指标做成一张「管道健康看板」显示每天的事件量、ID 绑定率、标签覆盖率、ETL 任务失败次数。这张看板不需要复杂工具QuickSight 连一张每天 ETL 结束时写入的 summary 表就够。我习惯在每天上班前打开这张看板先看数据量有没有异常波动再看绑定率是否稳定最后看各张表的标签覆盖有没有新缺口。这件事养成了习惯之后CDP 基本不会突然给你「惊喜」。我现在接手任何 CDP 项目都会在第一天先把这套质量基线搭好再谈功能迭代。方向本身是能复现的边界和成本控制得好它能成为团队的数据底座希望帮到你。本文还有配套的精品资源点击获取
返回列表