ARTICLE DETAIL

资讯详情

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

基于Hadoop+Spark的汽车销售数据分析平台与数仓实践

基于Hadoop+Spark的汽车销售数据分析平台与数仓实践 1. 项目整体设计思路1.1 为什么偏偏是这套技术栈汽车销售数据分析和普通的小规模数据分析不太一样。单店或者单一区域的销量数据用Excel、MySQL就能处理撑死上亿行也就那几百兆。但一旦数据来自多品牌、多门店、多渠道按天、按品牌、按车型、按经销商粒度累加日积月累就是几十GB甚至TB级别。数据量大只是一方面更麻烦的是分析口径经常调整今天要看厂家批发和终端零售的差距明天要看各车型的价格段渗透率后天又要看客户增换购的行为轨迹。传统数据库在这种情况下建索引、写复杂SQL都会非常痛苦。所以这套系统的选型思路很明确用HDFS作为海量原始数据的统一存储底座用Hive把原始数据管起来并做成数仓分层用Spark处理复杂ETL和高频聚合计算再用Python做数据采集脚本和可视化后端接口。四者不是堆砌技术而是各管一段Hadoop管存、Hive管仓、Spark管算、Python管采集和展示。这套组合最大的好处是每一层都有人干专门的活替换起来也灵活。比如后期想把Hive替换成Spark SQL直接跑或者把可视化从前端框架换成Superset都不用推翻整体架构。顺便说一句这个项目对个人学习也特别友好。Hadoop、Spark、Hive是招聘市场上大数据岗位问得最多的三个组件Python是数据分析的主流工具。把四者串成一个完整系统理解程度远高于只看理论或者只跑官方Demo对面试和工作都有直接帮助。1.2 整体架构与数据流转路径系统的数据链路大致如下数据源层业务系统导出的汽车销售订单明细、客户档案、经销商资料、车型配置信息以CSV/JSON格式落地。采集层Python脚本定时扫描文件目录解析数据后写入HDFS指定目录按日期分区存放。存储层HDFS按/data/raw/sales/dt2024-05-xx的目录结构保存源头数据。数仓层Hive建立外部表挂载原始文件再通过Spark SQL/Spark Core做清洗生成DWD明细层和ADS应用层结果表。分析层Spark读取Hive分析层表跑销售排行、趋势、份额、客户画像等指标回写Hive结果表。可视化层Python Flask提供HTTP接口读取结果表数据并返回JSON前端用ECharts绘制图表。你可能会问为什么不用Flume或者Sqoop去做采集Flume更适合日志流式采集Sqoop偏向关系型数据库和HDFS互导。这个场景里汽车销售数据是业务库定期导出的批量和半结构化文件用Python写脚本最直接解析逻辑可以灵活加各种规则开发调试也快。这也是大多数中小团队的实际做法。1.3 业务指标与数据口径定义起步阶段最重要的事情不是敲代码而是把业务流程和数据口径聊清楚。汽车销售领域通常涉及两个核心口径批发销量厂家卖给经销商和终端零售销量经销商卖给客户。做可视化之前必须明确你的统计口径否则图表做得再好看业务方一问就露馅。我这个项目以终端零售为主核心字段设计如下订单维度订单编号、成交日期、门店ID、销售顾问ID。车辆维度品牌、车型、厂商指导价、成交价、车身颜色、排量、能源类型。客户维度客户ID、性别、年龄段、所在城市。渠道维度销售渠道类型4S店、直营店、线上订单。指标体系我分了三类销量与规模指标日销量、月销量、累计销量、同比环比。结构与排名指标品牌份额、车型TOP10、价格区间分布、区域销量排行。质量与效率指标平均成交折扣率、新能源渗透率、成交周期、销售顾问人均销量。这些指标会贯穿后面的建表、分析和可视化全过程。提前把它们定义好后面所有工作都围绕它们展开不会跑偏。2. 环境准备与集群搭建2.1 Hadoop部署要点从单机到集群的取舍很多初学者会纠结要不要直接搭三节点集群我建议分两种情况。如果你本机内存小于16GB老老实实先搭Hadoop伪分布式也就是单节点上同时跑NameNode、DataNode、ResourceManager和NodeManager。伪分布式模式下的组件行为和生产集群基本一致跑通流程后再横向扩展成集群也不难。如果内存够大、机器够多就直接上三节点或五节点集群。Hadoop版本我选的是3.3.x。相比2.x版本3.x默认支持了Java 8以上生态更完整NameNode的联邦机制也更成熟。这里有一个很关键的配置点要提一下伪分布式模式下core-site.xml 里 fs.defaultFS 设置为 hdfs://localhost:9000hdfs-site.xml 里 dfs.replication 设置为1。集群模式下 replication 至少要设置为2。另外一个容易被忽略的细节是 ssh 免密登录。集群模式下启动脚本要跨节点拉起进程没有配置免密启动的时候会反复让你输密码非常影响体验。配置完记得用ssh localhost验证一次确保不用密码能登录。还有一点端口冲突是个高频问题。项目里如果已经装了别的服务占用了8088YARN Web UI或9870NameNode Web UI就得在 yarn-site.xml 和 core-site.xml 里换端口。我碰到过一次服务怎么都起不来排查半天发现是端口被占用启动日志里其实已经写了只是没仔细看。2.2 在Hadoop之上部署Hive并整合SparkHive的定位很明确它是跑在Hadoop上的数据仓库工具负责把SQL翻译成MapReduce或Spark作业。我在项目里用Hive主要做两件事一是把HDFS上的原始文件建表管理起来提供统一的SQL查询入口二是作为Spark的数据源让Spark能直接读写Hive表。Hive部署前必须先准备好MySQL作为元数据库因为Hive默认自带的Derby不支持多会话并发项目一跑就会被锁死。在 hive-site.xml 里配置好数据库连接地址比如jdbc:mysql://localhost:3306/hive_metastore然后执行schematool -initSchema -dbType mysql初始化元数据。要把Spark整合进来核心是两处配置!-- hive-site.xml -- property namehive.execution.engine/name valuespark/value /property property namespark.master/name valueyarn/value /property另外还需要把Spark的jar包关联到Hive的lib里面比较快的做法是在 HADOOP_CLASSPATH 加上Spark相关路径或者直接做软链接。整合完可以跑一个简单的select count(*) from test_table验证作业引擎是否切换到了Spark正常的话YARN的ResourceManager页面上能看到Spark的Application。2.3 Python环境与连接组件配置Python在本项目里负责三块数据采集脚本、Flask可视化接口、数据分析辅助脚本。版本我推荐Python 3.8太新的版本和某些Hadoop相关组件可能存在兼容性问题。要连HDFS读写文件用hdfs库即可不需要装复杂的大数据客户端pip install hdfs pymysql flask pandas pyhive其中hdfs库通过WebHDFS协议访问HDFS需要在 core-site.xml 里把dfs.webhdfs.enabled设置为true。pyhive则用于Python直接查询Hive表。在Windows上开发、Linux上跑任务容易出现编码问题我强烈建议所有文本统一UTF-8。尤其是Python脚本里面处理中文数据时文件头部最好声明# -*- coding: utf-8 -*-拼接HDFS路径时尽量不要用中文目录名避免编码不一致导致找不到文件。3. 数据采集与数据质量保障3.1 模拟数据生成让项目跑起来的生命线真实汽车销售数据在没有业务方配合的情况下很难拿到但项目要跑通必须有数据。所以第一步最好自己动手写一个模拟数据生成器把真实业务字段和分布特征模拟出来。这不是造假而是用符合业务规律的数据来验证整个系统。模拟数据要考虑三点数据量要能撑起集群建议单日生成不少于20万条订单记录字段要符合前面提到的指标体系分布要有规律比如新能源品牌的销量要随时间呈现上升趋势传统燃油车品牌保持相对平稳某些价格区间不能全部集中在一个品牌上。我写的生成脚本结构大概是这样import random import datetime import csv brands [BYD, Tesla, Toyota, VW, Benz, BMW, Audi, Honda] cities [北京, 上海, 广州, 深圳, 成都, 杭州, 武汉, 西安] def generate_daily_orders(date_str, target_count200000): rows [] for i in range(target_count): brand random.choice(brands) model f{brand}-{random.randint(1, 8)} price round(random.uniform(8, 55), 2) deal_price round(price * random.uniform(0.88, 1.02), 2) rows.append([ datetime.datetime.now().strftime(%Y%m%d%H%M%S) str(i).zfill(6), date_str, brand, model, price, deal_price, random.choice(cities), random.choice([4S店, 直营店, 线上]), random.choice([男, 女]), random.randint(22, 58) ]) return rows生成后按天保存成一个CSV文件文件名带上日期方便后续分区。重点提醒CSV里所有字段都加双引号避免字段里出现逗号把列弄错位。汽车品牌名称是英文还行如果是中文品牌还要额外注意UTF-8的编码输出。3.2 Python采集脚本写入HDFS采集脚本的核心逻辑是扫描某个本地目录把当天的新文件经过简单校验后通过WebHDFS上传到HDFS。需要注意HDFS上的目录设计要按分区来即/data/raw/sales_orders/dt2024-05-xx/orders.csv。这样后面建立Hive外部表时可以直接用分区目录挂载不用再写临时表做数据搬运。WebHDFS上传代码from hdfs import InsecureClient client InsecureClient(http://localhost:9870, userhadoop) client.makedirs(/data/raw/sales_orders/dt2024-05-20) with open(sales_orders_20240520.csv, rb) as f: client.write( /data/raw/sales_orders/dt2024-05-20/orders.csv, dataf, overwriteTrue )上传之前最好做一次基础校验空文件过滤、行数统计、表头字段数量检查。空文件上传会污染Hive表导致count时出现奇怪结果。另外一个坑是WebHDFS写文件默认会压缩或者分块如果文件名带中文建议统一重命名成ASCII字符省去很多麻烦。3.3 数据清洗与质量检查原始数据进了HDFS不代表就能直接分析脏数据是分析系统最大的敌人。我对这个项目里的清洗逻辑做了这些处理空值处理成交价为空或者为0的记录直接剔除因为做价格区间分析时这种数据没有意义。异常值处理成交价高于指导价的折扣比大于1.05或者低于指导价5折的记录标记为异常单独落一张异常数据表不参与聚合。重复值处理订单编号是天然主键用Distinct的方法去掉重复订单号。业务规则校验城市维度必须存在于城市字典表中不存在则归入“未知区域”。清洗可以用Spark DataFrame算子完成。写Spark代码时有一个调优细节要特别注意读CSV时如果字段类型推断错误会严重影响后续分析。所以每次用spark.read.csv都要手动指定schema而不是依赖inferSchema。指定Schema能显著提升读取速度也能避免因为某些脏数据导致推断类型不稳定的问题。from pyspark.sql.types import StructType, StructField, StringType, DoubleType, IntegerType schema StructType([ StructField(order_id, StringType(), True), StructField(order_date, StringType(), True), StructField(brand, StringType(), True), StructField(model, StringType(), True), StructField(guide_price, DoubleType(), True), StructField(deal_price, DoubleType(), True), StructField(city, StringType(), True), StructField(channel, StringType(), True), StructField(gender, StringType(), True), StructField(age, IntegerType(), True) ])4. Hive建仓与Spark数据分析实战4.1 数仓分层设计汽车销售数据的分析场景很多直接在原始表上写SQL会越写越乱。所以参照数仓分层的标准做法拆成三层ODS层保持原始数据不变建外部表字段和源文件保持一致。这一层的作用是还原现场出了问题还可以回溯检查。 DWD层对ODS做清洗、去重、规范化。比如把订单时间和分区时间拆开把不规范的渠道字段映射成标准的枚举值把城市地区补充上大区信息。 ADS层面向业务结果沉淀应用表比如日销量汇总表、品牌月度排名表、价格区间分布表。这一层的数据量不大但每张表都对应一个可视化图表的指标。建表时有一个关键点ODS和DWD层用外部表ADS层用内部表。原因是ODS/DWD对应的是HDFS上的原始文件和清洗中间结果一旦任务重跑只需要改HDFS文件路径即可而ADS是结果数据生命周期跟随Hive管理删表就删数据方便彻底更新。4.2 Hive分区表与关键DDL原始表一定要做分区分区字段用日期。这样做的好处有两方面查询时能通过分区裁剪只扫描需要的目录日常重跑某一批次数据时不会影响其他日期的数据。建表语句示例CREATE EXTERNAL TABLE dwd_sales_clean ( order_id STRING, order_date STRING, brand STRING, model STRING, guide_price DOUBLE, deal_price DOUBLE, city STRING, channel STRING, gender STRING, age INT, region STRING ) PARTITIONED BY (dt STRING) ROW FORMAT DELIMITED FIELDS TERMINATED BY , STORED AS TEXTFILE LOCATION /warehouse/dwd/sales_clean;建表后需要手动添加分区或者用修复命令让Hive自动识别HDFS上的新分区目录MSCK REPAIR TABLE dwd_sales_clean;在汽车销售场景里时间字段有一个很典型的坑订单日期是字符串类型例如 “2024-05-20”而分区字段 dt 也是字符串类型。写SQL时最容易犯的错误是直接把两个字段比较导致日期过滤条件完全不生效。要用to_date(order_date)转换成统一日期格式再过滤。4.3 Spark读取Hive表进行分析数据进入DWD层后分析任务交给Spark。Spark有两种常用方式对接Hive一种是在Spark SQL里直接写HQL操作Hive表另一种是SparkSession加载Hive的表数据生成 DataFrame然后按业务逻辑做数据处理。第一种方式简单适合临时查数。第二种方式适合复杂分析和ETL任务可控性更强。from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(AutoSalesAnalysis) \ .enableHiveSupport() \ .config(spark.sql.shuffle.partitions, 20) \ .getOrCreate() df spark.sql(SELECT * FROM dwd_sales_clean WHERE dt2024-05-20)用Spark做聚合计算的一大优势是分布式。但要注意spark.sql.shuffle.partitions这个参数直接影响Shuffle后的分区数量。默认设置是200在小规模数据上反而会导致大量小任务和文件碎片我建议根据数据量调整到20到50之间能明显减少小文件产生。4.4 核心分析SQL示例销量排行与窗口函数日销量趋势和品牌排行是最常用的分析场景。日销量趋势SQL很简单按天group by再求和就行。品牌排行则要用到窗口函数这也是Hive面试里的高频考点。计算每个品牌在当月的销量排名SELECT brand, month, sales_cnt, rn FROM ( SELECT brand, month, SUM(deal_cnt) AS sales_cnt, ROW_NUMBER() OVER (PARTITION BY month ORDER BY SUM(deal_cnt) DESC) AS rn FROM dwd_sales_clean GROUP BY brand, month ) t WHERE rn 10;窗口函数ROW_NUMBER() OVER (PARTITION BY month ORDER BY sales_cnt DESC)是在分组内排序生成序号。这里要特别强调一下执行顺序先GROUP BY后窗口计算。如果用WHERE rn 10直接过滤在Hive的旧版本里是不支持在SELECT产生的别名上直接过滤的所以要包一层子查询。我在这个项目里踩过这个坑需要记住窗口函数的过滤必须在外面包一层。新能源渗透率这个指标在汽车销售分析里也非常重要定义为新能源车型销量占总销量的比例SELECT dt, SUM(CASE WHEN energy_type 新能源 THEN 1 ELSE 0 END) / COUNT(1) AS penetration_rate FROM dwd_sales_clean GROUP BY dt;一个容易忽略的点如果能源类型字段里有NULL值渗透率会偏低。所以清洗时就要把缺失的能源类型统一填充为“未知”并单独标注否则比率计算会被污染。4.5 客户画像与价格区间分析客户画像分析对汽车行业特别有价值。经销商不仅要知道卖了多少车更要知道谁在买车、偏好什么价位区间的车。价格区间分析的做法是把成交价分成几个区间经济型10万以下、中端10-20万、中高端20-35万、高端35万以上然后再看每个区间的销量和占比。SELECT CASE WHEN deal_price 10 THEN 经济型 WHEN deal_price 20 THEN 中端 WHEN deal_price 35 THEN 中高端 ELSE 高端 END AS price_level, COUNT(1) AS sales_cnt, ROUND(COUNT(1) * 100.0 / SUM(COUNT(1)) OVER(), 2) AS sales_ratio FROM dwd_sales_clean GROUP BY CASE WHEN deal_price 10 THEN 经济型 WHEN deal_price 20 THEN 中端 WHEN deal_price 35 THEN 中高端 ELSE 高端 END;如果只看均价会掩盖不同价格带之间的差异。这里用了SUM(COUNT(1)) OVER()做窗口聚合可以在一个SQL里同时看到各区间销量和总销量占比不用额外写子查询。这个方法是分析类SQL里比较实用的技巧能减少一次作业运行时间。5. 可视化报表搭建5.1 可视化选型与整体方案可视化层我建议采用Python Flask ECharts的组合。Flask负责提供JSON接口ECharts在前端渲染图表。这套组合的好处是轻量、开发速度快、组件丰富而且完全不需要额外买商业工具。把结果表数据转换成JSON前端收到后直接喂给ECharts即可。如果你的团队更倾向于开箱即用也可以考虑Superset。Superset是Airbnb开源的可视化工具原生支持Hive、Spark SQL等数据源配好连接之后直接在界面上拖拽出图表。但Superset的自定义程度和交互精细度不如ECharts灵活。个人项目或者小团队我推荐Flask ECharts团队协作和需要自助分析的场景选Superset更省心。5.2 Flask后端接口设计与实现后端接口的设计思路很简单每个接口对应一个ADS层结果表或一个固定的SQL查询查询结果转成JSON格式返回前端。from flask import Flask, jsonify from pyhive import hive app Flask(__name__) def query_hive(sql): conn hive.Connection(hostlocalhost, port10000, usernamehadoop) cursor conn.cursor() cursor.execute(sql) columns [desc[0] for desc in cursor.description] rows cursor.fetchall() result [dict(zip(columns, row)) for row in rows] cursor.close() conn.close() return result app.route(/api/daily_sales) def daily_sales(): sql SELECT dt, SUM(sales_cnt) AS total_sales FROM ads_daily_sales GROUP BY dt ORDER BY dt data query_hive(sql) return jsonify({code: 0, data: data})有一个实战细节PyHive连接Hive默认走HiveServer2在启动Hive前需要在配置文件里启用HiveServer2服务。没启用的话接口调用直接报连接拒绝。接口返回的数据量也要控制比如品牌排行只传给前端Top10不要一次传上万行否则前端渲染卡顿、接口响应慢整个仪表盘体验很差。5.3 核心图表类型与使用场景不同分析场景适合不同图表类型。我的实践结果如下日销量趋势折线图横轴日期纵轴销量可以叠加7日移动平均线。品牌销量TOP10横向柱状图品牌名称较长时横向排布更易读。品牌市场份额饼图或环形图展示各品牌占比。价格区间分布堆叠柱状图同时展示各价格段内品牌分布。区域销量地图地图热力图表现各城市销量密度。新能源渗透率双轴折线图一条线是销量一条线是渗透率时间相同对比趋势。图表颜色不建议使用大量高饱和色汽车销售场景更偏向沉稳、大气的配色。如果数据有波动图表上的数值标签会显得拥挤可以通过formatter函数把数值转换成“万”“亿”等单位再显示。5.4 性能优化结果物化与预聚合可视化查询最忌讳的是用户每次点击都触发一次全量聚合。Spark跑一次全量聚合可能要几分钟前端用户不可能等。解决办法是把分析结果预先算好写入ADS层结果表。每天定时调度跑一次任务把当天的汇总指标算完可视化接口只需要查结果表毫秒级返回。另外ADS层的数据量虽然不大但依然建议在Hive表上做分区和排序或者直接构建成Parquet格式的表提升查询性能。ADS结果表可以按月归档查询时指定时间范围减少扫描量。对于实时性要求高的场景可以做增量计算每天只跑当天数据再和前一天结果做合并。不过做增量合并时要注意重复数据和漏算数据的问题宁可把所有日期分区重算也要保证结果准确。数据正确性永远优先于效率。6. 常见问题与排查经验6.1 Hive与Spark元数据不一致症状用Spark SQL查询Hive表能看到表结构但查不到数据或者用Hive能查到数据用Spark查是空结果。原因分析Spark和Hive共用同一个Metastore时需要配置相同的hive-site.xml。如果Spark的配置里没有读取到Hive的Metastore地址Spark就会使用自己内置的Derby数据库元数据和Hive完全隔离自然看不到数据。解决办法将Hive配置目录下的 hive-site.xml 复制到Spark的 conf 目录下重启Spark相关服务。然后验证一下是否生效spark-sql --execute SHOW TABLES;能正常显示所有Hive表就代表关联成功。6.2 数据倾斜导致任务卡死症状某个Spark作业整体进度卡在99%不动最后失败查看Stage看到某几个Task处理的数据量是其他Task的几十倍。原因分析数据倾斜通常由GROUP BY某个键值分布不均导致。比如品牌字段中可能有某个品牌的数据量远大于其他品牌。那么多条记录分发到同一个Reduce任务就会拖垮任务进度。解决办法可以先定位倾斜键然后用加盐法。给倾斜的键加随机前缀先分散聚合一次去掉前缀后再聚合一次。还有一种更简单的做法过滤掉极端倾斜的键值单独处理后再合并结果。连续几天观察到某个品牌订单特别多就可以把这个品牌抽出来单独跑分析避免影响整体任务。6.3 中文乱码问题症状Hive表查询结果里中文品牌、城市字段显示成???或\u0000。原因分析Hive默认的字符集和文件编码不一致。文件在写入HDFS时是UTF-8但Hive表的SERDE可能按Latin-1解析导致中文乱码。还有一种情况是Spark在写入结果表时spark.sql.session.timeZone设置不当导致时间类字段错乱和中文无关但容易被混淆。解决办法建表时指定ROW FORMAT DELIMITED FIELDS TERMINATED BY , COLLECTION ITEMS TERMINATED BY \n并确保文件编码为UTF-8且无BOM。写CSV文件时不要加BOM头因为Hive读取时BOM会被当成字段内容导致第一个字段总是出现不可见字符。6.4 小文件过多影响整体性能症状HDFS NameNode内存占用飙升Spark读取数据极慢每次扫描都需要打开大量文件。原因分析生产过程中按天分区、多任务写入容易产生大量小文件。尤其Spark任务默认每个分区写一个文件如果输出分区数设得太大会生成大量MB级别的小文件。解决办法调整spark.sql.shuffle.partitions避免分区数过大任务完成后对ODS/DWD层做一次文件合并INSERT OVERWRITE TABLE dwd_sales_clean PARTITION (dt2024-05-20) SELECT * FROM dwd_sales_clean WHERE dt2024-05-20 DISTRIBUTE BY dt;用DISTRIBUTE BY让相同日期的数据写进同一个分区文件合并后再查询扫描速度明显提升。6.5 常见问题速查表问题现象可能原因快速解决办法Spark作业一直PendingYARN资源不足调大yarn.nodemanager.resource.memory-mb或减少spark.executor.memoryHive查询一直卡在Map阶段数据文件有坏块或小文件过多在HDFS上检查文件大小分布先合并文件再查询PyHive连接超时HiveServer2未启动或端口被占用启动HiveServer2并确认10000端口监听Python写入HDFS报错WebHDFS未开启修改hdfs-site.xml设置dfs.webhdfs.enabledtrue数据量翻倍任务重复执行使用INSERT OVERWRITE而非INSERT INTO覆盖分区6.6 排障方法论小结排障顺序很重要。先确认底层存储是否正常再查计算引擎日志最后才查应用代码。我在项目里养成一个习惯每次跑任务前确认三件事——HDFS空间充足、YARN资源池正常、Hive表存在且分区正确。这三件事任何一个出问题任务都会以各种诡异姿势失败。日志里报错信息80%都指向真正原因多半是配置没对齐或者路径写错真正意义上的逻辑Bug反而不多。7. 项目复盘与实操心得7.1 目录规划与命名规范项目跑起来之后文件目录和命名规范最容易乱。我踩过这个坑后把所有目录严格按功能划分/opt/autosales/ ├── collector/ # Python采集脚本 ├── etl/ # Spark清洗分析脚本 ├── sql/ # Hive建表与查询SQL ├── web/ # Flask可视化后端 └── docs/ # 业务文档与口径说明数据代码文件和SQL脚本分目录存放方便后续维护。脚本和SQL文件头部加上日期和版本注释改过一次就更新一次避免后面找不到哪个版本是最新的。7.2 调度方案建议这个项目里采集每天固定时间执行分析任务在采集之后跑可视化接口实时查询结果表。自动化利器推荐Apache Airflow没有条件引新组件的话退而求其次用Crontab也能兜底。Airflow可以管理依赖关系、失败重跑和日志追踪Crontab方式简单直接配合Shell脚本里加的“任务失败自动发邮件”功能也能达到基本效果。调度时考虑两点各个任务之间有先后依赖比如必须先采集、后清洗、再分析失败任务要能重跑且不能重复写入数据。在Spark写结果时用INSERT OVERWRITE而不是INSERT INTO是规避重复数据最简单的方法。7.3 业务理解与技术实现的平衡做这类系统时间长了最深的体会是技术本身难度不大真正难的是把业务理解转化为技术实现。汽车销售分析涉及经销商、厂家、客户三方视角同样一张销量表给管理层看要突出总销量和增长率给销售团队看要细化到车型和门店给市场部看要聚焦品牌份额和价格段。在做可视化之前先确认好“这张图的受众是谁”比把数据做精确更重要。另外指标口径一定要有文档记录。比如新能源渗透率分母是所有销量还是新能源加燃油车总销量各个人理解不同。没有统一口径换个人来维护项目图表数据就对不上。我在项目里用docs目录维护了一份指标口径说明每次开会讨论都以此为准。7.4 最后分享一个使用技巧Spark读取Hive表时如果只查几天的分区数据建议把spark.sql.hive.convertMetastoreParquet设为false避免Spark在读取时额外做元数据转换能节省一些启动时间。同时把常用查询的SQL保存成视图后续Spark脚本里直接引用视图名不需要反复拼接长SQL。跑完任务记得清理Spark临时文件目录和日志文件磁盘会被这些隐形的临时文件慢慢吃掉。HDFS的回收站默认保留48小时删掉的数据还能找回磁盘不够时先清空回收站再报错。这套系统从零到跑通真正花费时间最多的不是写代码而是等任务跑、找配置错误、梳理数据口径。把链路跑通了之后加新车型、新指标、新门店都只是加一条SQL的事。所以不要急着写代码先把架构和数据流程想清楚后面一切都会顺利很多。
返回列表