Spark酒店经营数据分析:间夜口径、数仓分层与倾斜调优 📅 发布时间:2026/9/18 3:41:44 👁 浏览次数: 酒店数据分析这件事很多人第一反应是不就是几张表 join 一下出几个指标。我一开始也这么想直到接手一家连锁酒店集团的月报改造项目才发现酒店行业的指标口径比电商复杂得多——间夜这个词能把一半的开发者绕进去。那次任务的核心就是把原来跑在 Excel 和单机 Python 上的酒店经营分析整条链路迁到 Spark 上每天凌晨定时跑出一份覆盖入住率、ADR、RevPAR、复购率、渠道分布的宽表供 BI 和运营直接取数。下面我把这套东西从数据摸底到任务提交的完整过程拆开讲包括踩过的坑和最后稳定下来的配置。适合已经会写 Spark 代码、但没做过酒店行业数据、或者做出来结果总被业务挑战的同行参考。1. 先把五张表的关系捋清楚再谈怎么写代码1.1 订单表是事实其余四张都是维度酒店数据模型的骨架其实很清晰一张订单事实表配酒店、房型、用户、渠道/城市四张维度表。订单表里存的是某人在某天订了某家酒店的某个房型从哪天住到哪天花了多少钱现在处于什么状态。酒店表提供房量、星级、所属品牌和城市房型表提供床型、可住人数、挂牌价用户表提供注册渠道和会员等级城市/渠道表提供地域和流量来源的归属。这套结构听起来和电商订单没区别但真正的差异藏在粒度上。电商的一条订单对应一个 SKU分析时基本不用考虑时间跨度。酒店订单是带时间区间的一条订单可能横跨 1 到 30 天而且每天占用的房间数还取决于订了几间。所以做酒店分析事实表的粒度最终必须落到间夜这一层也就是把一条订单住 3 晚 2 间拆成 6 个间夜记录。这个下沉动作如果没有在 DWD 层做完后面所有指标都会算错而且是那种看起来数值合理、实际全错的错法。我在项目里定的分层是ODS 保留原样落 Hive 分区DWD 做清洗和粒度下沉DWS 按酒店 × 日期和用户 × 日期两个主题聚合ADS 直接出报表口径。分层的意义在于口径调整只需要改 DWS不用动 ODS 的同步逻辑这在指标被业务反复挑战的时候能省大量时间。1.2 字段里埋着的三个坑第一个坑是金额单位。订单表的pay_amount存的是分而房型表的base_price存的是元。这两个字段如果直接拿来做 ADR平均房价结果会差 100 倍而且因为酒店房价量级本来就有几百到几千的跨度看数字的时候很容易自我说服这个价格好像也合理。我的做法是在 DWD 层统一转成元用decimal(18,2)而不是double避免累加时的浮点误差。第二个坑是状态码的语义。预付业务线和现付业务线的状态机根本不一样。预付有待支付→已支付→已入住→已完成现付压根没有已支付这个中间态多了个 noshow预订未到。如果用一个统一的status in (3,4)去筛有效订单现付订单里的 noshow 会被算进取消取消率直接虚高。解决办法是在 DWD 层加一个biz_type字段按业务类型分别映射状态。第三个坑是空值的语义。nights字段为 null 有两种可能一种是新数据没写进来另一种是历史遗留数据压根没这个字段。这两种情况的处理方式完全不同——前者要退回用离店日期减入住日期补算后者如果也这么补会把一批跨月长住订单算成正常间夜。我在清洗的时候加了个nights_source标记字段记录这个值到底是原始的还是推算的出问题的时候能一眼看出数据来源。提示拿到一份陌生的酒店数据先别急着写计算逻辑花半小时把每张表的字段空值率、状态码分布、金额的 min/max 都跑一遍。这一步能省掉后面至少两天的返工。2. ODS 到 DWD清洗规则不是越严越好2.1 状态机决定了哪些订单能进宽表先说一个反直觉的判断DWD 层不应该只保留有效订单。如果把取消订单、noshow 订单在清洗阶段就过滤掉那你永远算不出取消率也没法做预订到入住的转化漏斗。正确的做法是全部保留但给每条订单打上清晰的状态标签让下游按需筛选。我在项目里定义了四个布尔标记is_valid_stay有效入住、is_cancel取消、is_noshow未到、is_pending进行中。判定逻辑按biz_type分支走预付业务用已入住或已完成判有效现付业务用已完成判有效。这套标记写好之后DWS 层算指标就变成纯粹的条件求和逻辑一目了然也不会出现两个人算出两个数的情况。-- DWD 层状态标记Hive SQL 示意 SELECT order_id, user_id, hotel_id, room_type_id, channel_id, check_in_date, check_out_date, nights, room_cnt, CASE WHEN biz_type PREPAY AND order_status IN (3, 4) THEN TRUE WHEN biz_type POSTPAY AND order_status 4 THEN TRUE ELSE FALSE END AS is_valid_stay, order_status 5 AS is_cancel, order_status 7 AS is_noshow FROM ods.ods_hotel_order WHERE dt ${bizdate};这里有个细节值得展开room_cnt房间数字段经常是空的默认应该按 1 处理而不是 0。按 0 处理会让间夜数直接归零而且因为聚合用的是 sum个别订单归零在总量上完全看不出来。我在清洗时统一做了coalesce(room_cnt, 1)并在数据质量校验里加了一条room_cnt 为 0 的记录数必须等于 0。2.2 间夜数为什么不能直接用 datediff最直观的写法是datediff(check_out_date, check_in_date)这个算式在绝大多数情况下是对的但酒店行业有几个特殊场景会让它失效。第一个是钟点房入住和离店是同一天datediff 返回 0但实际产生了 1 个间夜的口径占用。第二个是跨零点入住的夜审逻辑部分酒店系统把凌晨 2 点前入住算作前一天的房check_in_date 里已经做过了日期回退如果你再按自然日期减就会少算一晚。我最终采用的口径是优先取原始nights为空时用 datediff 补但用greatest(..., 1)兜底保证任何一条有效入住的订单至少产生 1 个间夜。同时把钟点房单独用一个room_type_category标记出来报表里默认排除需要单独看的时候再放开。from pyspark.sql import functions as F base orders.filter(F.col(is_valid_stay)) \ .withColumn(nights_final, F.when(F.col(nights).isNotNull(), F.col(nights)) .otherwise(F.greatest( F.datediff(check_out_date, check_in_date), F.lit(1)))) \ .withColumn(room_cnt_final, F.coalesce(room_cnt, F.lit(1))) \ .withColumn(room_nights, F.col(nights_final) * F.col(room_cnt_final))粒度下沉那一步我用sequence函数把入住到离店的日期展开成数组再 explode这样每个间夜都带上所属的自然日。数据量会膨胀但膨胀倍数等于平均入住天数乘以房间数酒店行业的平均值大概在 1.3 到 1.5 之间完全在可控范围。2.3 金额单位统一与分区写入金额处理上我坚持一个原则进数仓时就归一别留给下游。所有金额字段在 DWD 层统一转成元、统一用 decimal、统一把负数退款取绝对值后单独标记方向。这三条看起来像强迫症但实际项目里因为金额口径不一致被业务打回来的次数比因为性能问题被运维找上门的次数多得多。分区写入的配置也必须在这里定型。增量任务写 Hive 时如果不开动态分区覆盖每次都会把整个分区目录重写一遍数据量上去之后 IO 会非常难看。spark.conf.set(spark.sql.sources.partitionOverwriteMode, dynamic) spark.conf.set(hive.exec.dynamic.partition.mode, nonstrict) dwd.write.mode(overwrite).format(parquet) \ .partitionBy(dt) \ .saveAsTable(dwd.dwd_hotel_order_detail)注意partitionOverwriteModedynamic只在INSERT OVERWRITE语义下生效用 DataFrame 的overwrite模式配合partitionBy时行为取决于具体版本稳妥做法是建好表之后统一走 SQL 的INSERT OVERWRITE。3. 指标口径酒店行业和电商最大的区别在间夜3.1 入住率、ADR、RevPAR 三者的换算关系这三个指标是酒店经营的命根子但它们之间的换算关系经常被写错。定义本身很干净入住率 OCC 等于已售间夜除以可售间夜ADR 等于客房收入除以已售间夜RevPAR 等于客房收入除以可售间夜。三者关系是 RevPAR ADR × OCC。记住这个等式任何时候算出来的三个数不满足它就说明口径出问题了。难点全在可售间夜这个分母上。它不是简单地拿酒店表里的room_count乘以统计天数因为酒店会有装修停业、部分楼层关闭、新店未开业、老店已停业这些情况。我最后建了一张日快照维表dim_hotel_room_stock每天记录每家酒店当天实际可售的房量由酒店运营系统每日同步。这张表看起来不起眼但它是入住率能不能被业务认可的关键——用静态房量算出来的入住率在淡旺季会被严重扭曲。-- DWS酒店日粒度经营指标 SELECT s.dt, s.hotel_id, SUM(f.room_nights) AS sold_room_nights, SUM(s.room_count) AS available_room_nights, SUM(f.room_fee) AS room_revenue, ROUND(SUM(f.room_nights) / SUM(s.room_count), 4) AS occ, ROUND(SUM(f.room_fee) / SUM(f.room_nights), 2) AS adr, ROUND(SUM(f.room_fee) / SUM(s.room_count), 2) AS revpar FROM dwd.dwd_hotel_room_night f JOIN dim.dim_hotel_room_stock s ON f.hotel_id s.hotel_id AND f.dt s.dt GROUP BY s.dt, s.hotel_id;客房收入的定义也要说清楚它指的是房费不含餐饮、SPA、会议室这些。订单表里如果是打包价需要按比例拆分或者用一个配置化的拆分系数表来算。我在项目里维护了一张房费占比表按产品类型给不同系数虽然粗糙但比拍脑袋拆准确得多。3.2 复购率的三种口径选错了结论就废了复购率是这次改造里被业务改过最多遍的指标。常见的三种口径差异极大口径定义适用场景典型数值自然周期口径统计月内下单次数 ≥ 2 的用户占比快速看趋势5% 到 12%滚动窗口口径首次离店后 90 天内再次入住衡量真实忠诚度20% 到 35%会员生命周期口径一年内至少两次入住会员运营评估15% 到 25%同样一批数据三种口径能差出三四倍。业务方说复购率太低的时候先问清楚他们指的是哪一个。我在 DWS 层干脆把三种都算出来并存报表上标注口径名称谁要用哪个自己取。看起来是偷懒实际上减少了大量的来回确认。滚动窗口口径的实现用窗口函数最自然按用户分区按入住日期排序用lag取上一次离店日期再算时间差是否落在 90 天内。WITH stay_seq AS ( SELECT user_id, order_id, hotel_id, check_in_date, check_out_date, LAG(check_out_date) OVER ( PARTITION BY user_id ORDER BY check_in_date ) AS prev_check_out FROM dwd.dwd_hotel_order_detail WHERE is_valid_stay TRUE ) SELECT COUNT(DISTINCT CASE WHEN DATEDIFF(check_in_date, prev_check_out) BETWEEN 1 AND 90 THEN user_id END) / COUNT(DISTINCT user_id) AS repurchase_rate_90d FROM stay_seq;这段 SQL 有个容易忽略的点prev_check_out为 null 的首单用户也要进分母。我在第一次写的时候用了内连接过滤导致分母只剩老客复购率算出来 60% 多被业务当场质疑。3.3 渠道归因与取消率的口径陷阱渠道分析里最容易出错的是归因方式。一笔订单可能来自搜索引擎的落地页经过比价平台跳转最后在自有渠道完成支付。到底算哪个渠道的功劳酒店行业主流做法是四段归因首触渠道、末触渠道、结算渠道、订单归属渠道四个字段都存下来。报表上默认用订单归属渠道因为这是和结算分成直接挂钩的口径做投放效果评估时才切到首触。取消率的口径陷阱更隐蔽。分母如果用全部下单量那些下了单没付款就自动关闭的订单会被算进去取消率虚高得离谱。行业里比较被认可的口径是取消率 取消单量 / (有效入住单量 取消单量 noshow 单量)也就是只保留走完了预订流程的订单。这个定义我在项目文档里写了三遍并且在 DWS 表上加了口径注释字段就是为了防止后面接手的人按自己的理解重算。4. 倾斜排查从 WebUI 上一个 Task 卡住说起4.1 酒店数据的倾斜往往来自渠道而不是酒店第一次跑全量任务的时候99% 的 task 在 2 分钟内结束剩下 3 个 task 卡了 40 分钟还没动。打开 Spark UI 一看那个 stage 的 Task Duration 最大值 2400 秒中位数 45 秒差了 50 多倍。这就是典型的数据倾斜。反查 key 分布之后发现倾斜不是来自酒店维度而是来自渠道维度SELECT channel_id, COUNT(*) AS cnt FROM dwd.dwd_hotel_order_detail GROUP BY channel_id ORDER BY cnt DESC LIMIT 20;结果很典型头部两三个第三方渠道吃掉了 80% 以上的订单量剩下几十个渠道分 20%。这种分布几乎是酒店行业的标配因为流量本来就是高度集中的。同理酒店维度上也会有倾斜连锁品牌旗下的大店单店订单量可能是小店的几百倍。还有个容易被忽略的倾斜源是null key。用户表关联不上的订单user_id为 null所有 null 在 join 时会被分到同一个 partition如果这类订单有几十万条就会形成一个巨大的单点。4.2 完整的排查链路我把这套排查流程固化成了一个清单每次遇到慢任务都按顺序走一遍第一步看 Stage 页的 Task Duration 分布最大和 P75 差距超过 10 倍就基本确认倾斜。第二步看 Shuffle Read Size 的最大值和平均值确认是数据量倾斜而不是计算逻辑慢。第三步定位到具体算子是groupBy还是join是哪个 stage 触发的 shuffle。第四步反查 key 分布用上面那段 SQL把 top 20 的 key 和占比都列出来。第五步判断倾斜是真倾斜还是假倾斜——有时候是 join 条件写错了比如带了or条件导致笛卡尔积有时候是 key 里混了空字符串和 null 两种空值看起来是两个 key实际是一个巨大的分组。第六步确认是否为数据异常导致比如某个月数据重复导入让某个 key 的量凭空翻倍。这一步我用分区行数环比做判断波动超过 30% 就先去查上游。4.3 加盐、广播、两阶段聚合怎么选解法不是越复杂越好我的选择顺序是能广播就广播不能广播再看能否拆表最后才加盐。小表关联大表的情况直接广播map 端 join 不打散一步解决。酒店维表、房型维表、城市维表都属于这类行数在几万以内广播阈值默认 10MB 完全够用。from pyspark.sql.functions import broadcast hotel_dim spark.table(dim.dim_hotel).select( hotel_id, hotel_name, star, city_id, brand_id, room_count) fact order_night.join(broadcast(hotel_dim), onhotel_id, howleft)大表关联大表又必须按大 key 聚合的情况加盐是标准解法。核心思路是给大 key 加随机前缀打散做一次局部聚合再去掉前缀做全局聚合。盐值的数量按倾斜程度定一般 16 到 64 之间太多会引入额外的 shuffle 开销。SALT 32 salted fact.withColumn(salt, (F.rand() * SALT).cast(int)) \ .withColumn(salted_key, F.concat_ws(_, F.col(channel_id), F.col(salt))) partial salted.groupBy(salted_key).agg( F.sum(room_nights).alias(rn), F.sum(room_fee).alias(fee)) final partial.withColumn(channel_id, F.split(salted_key, _)[0]) \ .groupBy(channel_id).agg( F.sum(rn).alias(room_nights), F.sum(fee).alias(room_fee))还有一种情况是少数几个 key 极端大加盐也救不了那就把热点 key 单独拆出来算剩下的走常规路径最后 union。这种写法代码丑但效果最直接。实测下来把 top 2 渠道单独处理之后那个 stage 的耗时从 40 分钟降到 6 分钟。经验Spark 3.x 的 AQE 已经能自动处理一部分倾斜开spark.sql.adaptive.skewJoin.enabledtrue之后超过阈值的 partition 会被自动拆成子 partition。但这不能完全替代手工处理尤其是groupBy引起的倾斜AQE 的覆盖能力有限。5. 提交到 YARN 之前参数得一条条算过5.1 executor 的核数与内存怎么配资源配置这件事没有万能公式但有个可靠的推导思路先看数据量再看单分区处理量最后反推 executor 数量。我这次任务的基础数据量是订单明细 30 亿行左右单日增量约 1200 万行全量刷新的 shuffle 数据量大概 500GB。按单 task 处理 128MB 到 256MB 的经验值需要的 task 数在 2000 到 4000 之间。executor 的配置上我选的是 4 核 8G 内存memoryOverhead 加 2G也就是每个 executor 实际占 10G 容器内存。核数不宜超过 5超过之后单 executor 的 GC 压力会明显上升磁盘 IO 也会成为瓶颈。spark-submit \ --master yarn \ --deploy-mode cluster \ --name hotel_dws_daily \ --queue root.prod.etl \ --driver-memory 4g \ --executor-memory 8g \ --executor-cores 4 \ --conf spark.executor.memoryOverhead2g \ --conf spark.dynamicAllocation.enabledtrue \ --conf spark.dynamicAllocation.minExecutors5 \ --conf spark.dynamicAllocation.maxExecutors60 \ --conf spark.dynamicAllocation.executorIdleTimeout120s \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.sql.adaptive.skewJoin.enabledtrue \ --conf spark.sql.adaptive.coalescePartitions.enabledtrue \ --conf spark.sql.shuffle.partitions400 \ --conf spark.sql.session.timeZoneAsia/Shanghai \ --conf spark.sql.sources.partitionOverwriteModedynamic \ hotel_dws_daily.py内存溢出的排查有个信号要记住如果日志里出现Container killed by YARN for exceeding memory limits那是物理内存超了要加 memoryOverhead如果是 Java 的OutOfMemoryError: Java heap space那是堆内存不够要加 executor-memory。这两种问题的解法完全不同别搞混。5.2 shuffle 分区数不是越大越好spark.sql.shuffle.partitions默认 200这个值在数据量大的任务里明显偏小会导致单个 partition 数据量过大、频繁 spill 到磁盘。但也别一上来就设成几千分区太多会让 task 调度开销和小文件问题同时爆发。我的估算方式是分区数取并发 task 数的 2 到 3 倍。这个任务的并发 task 数在 60 个 executor 乘 4 核等于 240所以分区数设 400 到 700 比较合适。开了 AQE 之后coalescePartitions会在运行时自动合并过小的分区所以初始值可以略微设大一点让 AQE 去收敛。另外spark.sql.adaptive.advisoryPartitionSizeInBytes默认是 64MB酒店数据的单行比较宽维表 join 之后字段多我调到了 128MB减少输出文件数。5.3 动态分区写入与小文件治理每天的全量任务会往 Hive 里写几百个分区如果不做治理一个月下来小文件能堆到几十万个元数据压力非常大。我在写出口加了两个动作一是用repartition按分区字段重分布控制每个分区的文件数二是开启 AQE 的分区合并让最终落地的文件大小落在 128MB 附近。target_cols [dt, city_id] result.repartition(200, *target_cols) \ .write.mode(overwrite).format(parquet) \ .partitionBy(*target_cols) \ .saveAsTable(dws.dws_hotel_daily_metrics)repartition的列选择有讲究要选基数适中的维度。dt的基数是 30city_id的基数是 300 左右两个组合起来大概 9000 个分区每个分区分到的文件数就能控制在个位数。如果选了个基数上千的字段反而会制造更多小文件。提示写出口不要用coalesce它会导致上游计算阶段并行度骤降长任务会明显变慢。要压文件数就用repartition代价是一次额外的 shuffle但可控。6. 任务跑完之后校验、增量与调度6.1 数据质量校验清单任务跑成功不等于数据可用。我在 DWS 表落地之后加了一层校验任务跑完再通知下游校验不通过就阻断。这套清单现在固定六条校验项规则处理动作主键唯一order_id 的 distinct 数等于总行数阻断金额非负room_fee 0 的行数为 0阻断间夜为正room_nights 0 的行数为 0阻断日期合法check_out_date check_in_date阻断分区行数波动环比波动在 ±30% 以内告警关键字段空值率酒店 ID 空值率 0.1%告警阻断和告警要分开。硬性错误必须阻断因为下游拿到的就是错数据波动类的先告警人工看一眼再决定不然节假日流量波动大天天阻断会把调度搞瘫。校验的实现用一张规则配置表驱动每条规则写一段 SQL引擎自动跑。这套东西一开始觉得是额外工作量但上线三个月里它拦下了两次上游数据重复导入和一次字段类型变更价值远超投入。6.2 全量与增量的边界全量任务的成本高不可能天天跑。我的方案是DWD 层按天增量DWS 层按日期分区增量聚合涉及用户生命周期的那几个指标复购率、会员活跃走 T-1 到 T-30 的滚动窗口每天重算最近 30 天的分区。这样既保证了指标准确性又把计算量控制在一个可接受的范围内。滚动窗口重算的关键是把窗口大小做成参数而不是硬编码在 SQL 里。我遇到过业务临时要从 90 天改成 180 天看效果的情况参数化的写法十分钟就能改完上线。WINDOW_DAYS 180 # 从配置表读取便于调整 start_dt biz_date - timedelta(daysWINDOW_DAYS) user_window spark.sql(f SELECT user_id, MIN(check_in_date) AS first_stay, MAX(check_in_date) AS last_stay, COUNT(DISTINCT order_id) AS order_cnt FROM dwd.dwd_hotel_order_detail WHERE is_valid_stay TRUE AND dt BETWEEN {start_dt} AND {biz_date} GROUP BY user_id )6.3 出问题时先看哪里任务在调度上跑起来之后总会遇到各种异常。我的排障顺序是固定的先看 Spark UI 里有没有失败的 stage 和报错堆栈再看 driver 日志里有没有数据质量校验的失败信息最后才去看上游分区有没有正常产出。顺序反过来的话很容易在上游没数据的情况下排查半天代码逻辑浪费时间。有一类问题特别隐蔽任务显示成功但结果表是空的。这通常是因为partitionOverwriteModedynamic配合空分区写入时的行为——上游没有数据任务照样成功退出只是写出了一个空分区。我在任务末尾加了一条断言检查输出表的行数是否大于 0为 0 就直接抛异常让调度标红。最后分享一个我在实际操作中的习惯把每次调优前后的关键指标记在一张表里包括 shuffle 数据量、stage 耗时、spill 大小、输出文件数。这些数字看起来琐碎但当你想判断某个参数调整是不是真的有效时凭感觉是不可靠的。我记录到第 20 次左右的时候已经能凭经验预判某个改动大概能带来多少收益这个直觉完全来自于那张记录表。酒店数据分析这套东西指标口径的复杂度远高于 Spark 本身的技术难度把口径对齐这件事做扎实比多调几个参数有价值得多。