从零搭建大数据数仓:Hadoop + Hive + DataX + PostgreSQL 全流程复现
目录
- 一、架构与节点规划
- 二、基础环境准备(所有节点)
- 三、Hadoop 完全分布式集群配置
- 四、PostgreSQL 元数据库配置
- 五、Hive 安装与配置
- 六、DataX 数据同步(ODS 层)
- 七、DWD 明细层数据清洗
- 八、Python 聚合脚本与 ADS/DWS 表生成
- 九、总结
- 附录:完整 Python 聚合脚本
一、架构与节点规划
本项目模拟图书数据仓库,从业务库 PostgreSQL 出发,经过 ODS、DWD 直到 DWS/ADS 层,形成完整的分析链路。
数据流向:
PostgreSQL(业务库 192.168.168.1)
↓ DataX (postgresqlreader → hdfswriter)
HDFS ODS 层 (外部表)
↓ DataX (hdfsreader → hdfswriter 清洗)
HDFS DWD 明细层 (外部表)
↓ Python 聚合脚本
HDFS DWS / ADS 指标层 (外部表)
节点信息:
| 主机名 | IP | 运行组件 |
|---|---|---|
| master | 192.168.168.100 | NameNode, ResourceManager, PostgreSQL, Hive, DataX, Python 聚合脚本 |
| node1 | 192.168.168.101 | DataNode, NodeManager |
| node2 | 192.168.168.102 | DataNode, NodeManager |
说明:业务数据库 PostgreSQL 位于另一台物理机
192.168.168.1,仅用于提供原始数据。
所有操作默认使用root用户,生产环境请适当调整权限。
二、基础环境准备(所有节点)
执行节点:master, node1, node2
执行路径:任意
2.1 配置静态 IP
查看网卡名(例如 ens33):
ip addr | grep ens
编辑对应配置文件:
vi /etc/sysconfig/network-scripts/ifcfg-ens33
根据节点添加(master 的 IPADDR 为 100,node1 为 101,node2 为 102,其余字段一致):
IPADDR=192.168.168.100 # 各节点改为自己的 IP
NETMASK=255.255.255.0
GATEWAY=192.168.168.2
DNS1=114.114.114.114
DNS2=8.8.8.8
验证:
systemctl restart network
ip addr | grep inet
ping www.baidu.com -c 3
2.2 主机名与映射
# master 上
hostnamectl set-hostname master
# node1 上
hostnamectl set-hostname node1
# node2 上
hostnamectl set-hostname node2
三台都编辑 /etc/hosts,添加:
192.168.168.100 master
192.168.168.101 node1
192.168.168.102 node2
2.3 关闭防火墙与 SELinux
systemctl stop firewalld
systemctl disable firewalld
setenforce 0
vi /etc/selinux/config
# 修改 SELINUX=enforcing 为 SELINUX=disabled
2.4 更换阿里云 YUM 源
curl -o /etc/yum.repos.d/CentOS-Base.repo http://mirrors.aliyun.com/repo/Centos-7.repo
yum install -y wget vim net-tools lrzsz
成功显示截图
三、Hadoop 完全分布式集群配置
3.1 安装 JDK 与 Hadoop
执行节点:所有节点
路径:/model/hadoop(统一管理)
yum install -y java-1.8.0-openjdk java-1.8.0-openjdk-devel
# 创建软链接(根据实际路径调整)
ln -s /usr/lib/jvm/java-1.8.0-openjdk-1.8.0.412.b08-1.el7_9.x86_64 /usr/lib/jvm/jdk1.8
配置环境变量(所有节点):
echo "export JAVA_HOME=/usr/lib/jvm/jdk1.8" >> /etc/profile
echo "export PATH=\$JAVA_HOME/bin:\$PATH" >> /etc/profile
source /etc/profile
下载 Hadoop(仅 master 执行,随后分发):
mkdir -p /model/hadoop && cd /model/hadoop
wget https://mirrors.aliyun.com/apache/hadoop/common/hadoop-3.3.6/hadoop-3.3.6.tar.gz
tar -zxvf hadoop-3.3.6.tar.gz
mv hadoop-3.3.6 hadoop
所有节点设置 Hadoop 环境变量:
echo "export HADOOP_HOME=/model/hadoop/hadoop" >> /etc/profile
echo "export PATH=\$PATH:\$HADOOP_HOME/bin:\$HADOOP_HOME/sbin" >> /etc/profile
source /etc/profile
hadoop version # 验证
3.2 核心配置文件逐行注释
执行节点:master
路径:/model/hadoop/hadoop/etc/hadoop
以下所有配置文件均添加行内注释,解释每个属性的作用,并用 【替换】 标出需要根据实际情况修改的地方。
3.2.1 hadoop-env.sh
# 设置 Hadoop 使用的 Java 安装路径,必须指向实际 JDK 目录
sed -i 's|# export JAVA_HOME=|export JAVA_HOME=/usr/lib/jvm/jdk1.8|g' hadoop-env.sh
说明:hadoop-env.sh 是环境脚本,JAVA_HOME 指定 JDK 路径,所有 Hadoop 守护进程启动时都会读取。
3.2.2 core-site.xml
<?xml version="1.0" encoding="UTF-8"?>
<?xml-stylesheet type="text/xsl" href="configuration.xsl"?>
<configuration>
<property>
<name>fs.defaultFS</name>
<value>hdfs://master:9000</value>
</property>
<property>
<name>hadoop.tmp.dir</name>
<value>/model/hadoop/tmp</value>
</property>
<property>
<name>hadoop.proxyuser.root.hosts</name>
<value>*</value>
</property>
<property>
<name>hadoop.proxyuser.root.groups</name>
<value>*</value>
</property>
</configuration>
3.2.3 hdfs-site.xml
<?xml version="1.0" encoding="UTF-8"?>
<configuration>
<!-- 文件副本数,不能超过 DataNode 节点数量 -->
<property>
<name>dfs.replication</name>
<value>2</value> <!-- 【替换】如果只有1个 DataNode 则改为 1 -->
</property>
<!-- NameNode 元数据存储目录,建议使用多目录保证可靠性 -->
<property>
<name>dfs.namenode.name.dir</name>
<value>/model/hadoop/tmp/dfs/name</value>
</property>
<!-- DataNode 数据块存储目录 -->
<property>
<name>dfs.datanode.data.dir</name>
<value>/model/hadoop/tmp/dfs/data</value>
</property>
</configuration>
3.2.4 yarn-site.xml
<?xml version="1.0" encoding="UTF-8"?>
<configuration>
<!-- YARN 资源管理器所在主机名 -->
<property>
<name>yarn.resourcemanager.hostname</name>
<value>master</value> <!-- 【替换】改为你的 ResourceManager 主机名 -->
</property>
<!-- NodeManager 辅助服务,用于 MapReduce Shuffle -->
<property>
<name>yarn.nodemanager.aux-services</name>
<value>mapreduce_shuffle</value>
</property>
</configuration>
3.2.5 mapred-site.xml
<?xml version="1.0" encoding="UTF-8"?>
<configuration>
<!-- 指定 MapReduce 运行在 YARN 上 -->
<property>
<name>mapreduce.framework.name</name>
<value>yarn</value>
</property>
</configuration>
3.2.6 workers
master # 本行表示 master 节点也作为 DataNode 运行
node1 # 【替换】按实际节点名修改
node2
说明:此文件列出所有 DataNode(同时也作为 NodeManager),一行一个主机名。
3.3 分发配置与启动集群
执行节点:master
路径:/model/hadoop/hadoop/etc/hadoop
分发配置文件到 node1、node2:
for node in node1 node2; do
scp core-site.xml hdfs-site.xml yarn-site.xml mapred-site.xml workers hadoop-env.sh $node:/model/hadoop/hadoop/etc/hadoop/
done
格式化 NameNode(仅首次,二次执行会丢失数据):
hdfs namenode -format # 输入 Y 确认
启动 HDFS 和 YARN:
start-dfs.sh
start-yarn.sh
在各节点运行 jps 检查进程:
master
node1/node2
四、PostgreSQL 元数据库配置
执行节点:master
路径:/model/postgresql/data
yum install -y postgresql-server postgresql-contrib
mkdir -p /model/postgresql/data
chown -R postgres:postgres /model/postgresql
su - postgres -c "/usr/bin/initdb -D /model/postgresql/data"
修改监听地址(允许 Hive 连接):
vi /model/postgresql/data/postgresql.conf
# 将 #listen_addresses = 'localhost' 改为 listen_addresses = '*'
配置认证信任:
echo "host all all 127.0.0.1/32 trust" >> /model/postgresql/data/pg_hba.conf
echo "host all all ::1/128 trust" >> /model/postgresql/data/pg_hba.conf
创建 systemd 服务文件,指定数据目录:
mkdir -p /etc/systemd/system/postgresql.service.d
cat > /etc/systemd/system/postgresql.service.d/pgdata.conf << 'EOF'
[Service]
Environment=PGDATA=/model/postgresql/data
EOF
systemctl daemon-reload
systemctl start postgresql
systemctl enable postgresql
创建 Hive 元数据库及用户:
su - postgres -c "psql -c \"CREATE USER hiveuser WITH PASSWORD 'hive123';\""
su - postgres -c "psql -c \"CREATE DATABASE hive_metadata OWNER hiveuser;\""
su - postgres -c "psql -c \"GRANT ALL PRIVILEGES ON DATABASE hive_metadata TO hiveuser;\""
【替换】:用户名
hiveuser和密码hive123可根据需要修改,但必须与后续hive-site.xml保持一致。
五、Hive 安装与配置
执行节点:master
路径:/model
cd /model
wget https://repo.huaweicloud.com/apache/hive/hive-3.1.3/apache-hive-3.1.3-bin.tar.gz
tar -zxvf apache-hive-3.1.3-bin.tar.gz
mv apache-hive-3.1.3-bin hive
# 环境变量(所有节点)
echo "export HIVE_HOME=/model/hive" >> /etc/profile
echo "export PATH=\$PATH:\$HIVE_HOME/bin" >> /etc/profile
source /etc/profile
# 下载 PostgreSQL JDBC 驱动
cd /model/hive/lib
wget https://jdbc.postgresql.org/download/postgresql-42.6.0.jar
hive-site.xml(路径:/model/hive/conf/hive-site.xml)
用户名hiveuser,密码hive123
<?xml version="1.0" encoding="UTF-8"?>
<configuration>
<!-- ========== 元数据库配置(保留你的原有配置) ========== -->
<property>
<name>javax.jdo.option.ConnectionURL</name>
<value>jdbc:postgresql://localhost:5432/hive_metadata</value>
</property>
<property>
<name>javax.jdo.option.ConnectionDriverName</name>
<value>org.postgresql.Driver</value>
</property>
<property>
<name>javax.jdo.option.ConnectionUserName</name>
<value>hiveuser</value>
</property>
<property>
<name>javax.jdo.option.ConnectionPassword</name>
<value>hive123</value>
</property>
<property>
<name>hive.metastore.warehouse.dir</name>
<value>/user/hive/warehouse</value>
</property>
<property>
<name>hive.server2.thrift.bind.host</name>
<value>0.0.0.0</value>
<description>监听所有网卡,允许外部连接</description>
</property>
<property>
<name>hive.server2.thrift.port</name>
<value>10000</value>
<description>HiveServer2 端口</description>
</property>
</configuration>
初始化元数据模式:
schematool -initSchema -dbType postgres
# 看到 “schemaTool completed” 表示成功
验证:
hive -e "show databases;"
六、DataX 数据同步(ODS 层)
6.1 安装 DataX
执行节点:master
路径:/model
cd /model
wget https://datax-opensource.oss-cn-hangzhou.aliyuncs.com/202309/datax.tar.gz
tar -zxvf datax.tar.gz
echo "export DATAX_HOME=/model/datax" >> /etc/profile
echo "export PATH=\$PATH:\$DATAX_HOME/bin" >> /etc/profile
source /etc/profile
# 拷贝 PostgreSQL 驱动,使 DataX 支持 PG 读写
cp /model/hive/lib/postgresql-42.6.0.jar /model/datax/plugin/reader/postgresqlreader/libs/
cp /model/hive/lib/postgresql-42.6.0.jar /model/datax/plugin/writer/postgresqlwriter/libs/
# 自测
python /model/datax/bin/datax.py /model/datax/job/job.json
出现这个代表datax安装成功
6.2 pg2hdfs.json 配置详解
文件路径:/model/pg2hdfs.json
执行:python /model/datax/bin/datax.py /model/pg2hdfs.json
{
"job": {
"setting": {
"speed": { "channel": 1 } // 并发通道数,1 表示单线程同步
},
"content": [
{
"reader": {
"name": "postgresqlreader", // 读取插件:PostgreSQL
"parameter": {
"username": "postgres", // 【替换】业务库登录用户名
"password": "123456", // 【替换】业务库密码
"column": [ // 要读取的字段列表,按需调整
"id",
"book_name",
"author",
"press",
"publication_date",
"number_of_pages",
"price",
"\"ISBN\"", // 注意大小写敏感,含大写需加双引号转义
"score",
"number_of_reads"
],
"splitPk": "id", // 分片主键(当 channel>1 时用于并行切割)
"connection": [
{
"table": ["book_douban"], // 【替换】业务表名
"jdbcUrl": [
"jdbc:postgresql://192.168.168.1:5432/book_collection?connectTimeout=30"
// 【替换】IP 改为物理机真实 IP,book_collection 为数据库名
]
}
]
}
},
"writer": {
"name": "hdfswriter", // 写入插件:HDFS
"parameter": {
"defaultFS": "hdfs://master:9000", // 【替换】HDFS NameNode 地址
"fileType": "text", // 文件格式:文本
"path": "/user/hive/warehouse/ods.db/ods_book_douban", // ODS 目标路径
"fileName": "book_data", // 输出文件名前缀
"column": [ // 输出字段定义,类型可选 INT/STRING/LONG/DOUBLE 等
{"name": "id", "type": "INT"},
{"name": "book_name", "type": "STRING"},
{"name": "author", "type": "STRING"},
{"name": "press", "type": "STRING"},
{"name": "publication_date", "type": "STRING"},
{"name": "number_of_pages", "type": "STRING"},
{"name": "price", "type": "STRING"},
{"name": "ISBN", "type": "STRING"},
{"name": "score", "type": "STRING"},
{"name": "number_of_reads", "type": "STRING"}
],
"writeMode": "append", // 写入模式:追加(首次可用 nonConflict 或 truncate,按需)
"fieldDelimiter": "\t" // 列分隔符,需与 Hive 建表分隔符一致
}
}
}
]
}
}
运行前准备 HDFS 目录:
hdfs dfs -mkdir -p /user/hive/warehouse/ods.db/ods_book_douban
hdfs dfs -chmod -R 777 /user/hive/warehouse
执行同步:
python /model/datax/bin/datax.py /model/pg2hdfs.json
数据导入成功
6.3 创建 ODS 外部表
进入 Hive 客户端执行:
CREATE DATABASE IF NOT EXISTS ods;
USE ods;
CREATE EXTERNAL TABLE IF NOT EXISTS ods_book_douban(
id INT,
book_name STRING,
author STRING,
press STRING,
publication_date STRING,
number_of_pages STRING,
price STRING,
ISBN STRING,
score STRING,
number_of_reads STRING
)
ROW FORMAT DELIMITED
FIELDS TERMINATED BY '\t' -- 必须与 DataX 分隔符一致
LOCATION '/user/hive/warehouse/ods.db/ods_book_douban';
验证:
SELECT * FROM ods.ods_book_douban LIMIT 5;
七、DWD 明细层数据清洗
7.1 ods2dwd.json 配置详解
文件路径:/model/ods2dwd.json
执行:python /model/datax/bin/datax.py /model/ods2dwd.json
{
"job": {
"setting": {
"speed": { "channel": 1 }
},
"content": [
{
"reader": {
"name": "hdfsreader", // 读取插件:HDFS
"parameter": {
"path": "/user/hive/warehouse/ods.db/ods_book_douban/*", // ODS 路径,* 匹配所有文件
"defaultFS": "hdfs://master:9000", // 【替换】NameNode 地址
"fileType": "text",
"column": [ // 按列索引读取,index 从 0 开始
{"index": 0, "type": "long"},
{"index": 1, "type": "string"},
{"index": 2, "type": "string"},
{"index": 3, "type": "string"},
{"index": 4, "type": "string"},
{"index": 5, "type": "string"},
{"index": 6, "type": "string"},
{"index": 7, "type": "string"},
{"index": 8, "type": "string"},
{"index": 9, "type": "string"}
],
"fieldDelimiter": "\t",
"encoding": "UTF-8"
}
},
"writer": {
"name": "hdfswriter",
"parameter": {
"defaultFS": "hdfs://master:9000",
"fileType": "text",
"path": "/user/hive/warehouse/dwd.db/dwd_book_detail", // DWD 目标路径
"fileName": "dwd_book_detail",
"column": [ // 输出字段重命名,类型可按需转换
{"name": "id", "type": "BIGINT"},
{"name": "book_name", "type": "STRING"},
{"name": "author", "type": "STRING"},
{"name": "press", "type": "STRING"},
{"name": "pub_year", "type": "STRING"}, // 字段重命名:publication_date → pub_year
{"name": "pages", "type": "STRING"}, // number_of_pages → pages
{"name": "price", "type": "STRING"},
{"name": "isbn", "type": "STRING"}, // ISBN → isbn
{"name": "score", "type": "STRING"},
{"name": "read_count", "type": "STRING"} // number_of_reads → read_count
],
"writeMode": "truncate", // 每次运行前清空目标目录,保证幂等性
"fieldDelimiter": "\t"
}
}
}
]
}
}
准备 DWD 目录并执行:
hdfs dfs -mkdir -p /user/hive/warehouse/dwd.db/dwd_book_detail
hdfs dfs -rm -r /user/hive/warehouse/dwd.db/dwd_book_detail/* # 清空旧数据(如果存在)
python /model/datax/bin/datax.py /model/ods2dwd.json
清洗成功图(示例)
7.2 创建 DWD 外部表
CREATE DATABASE IF NOT EXISTS dwd;
USE dwd;
CREATE EXTERNAL TABLE dwd_book_detail (
id BIGINT,
book_name STRING,
author STRING,
press STRING,
pub_year BIGINT, -- 类型转为 BIGINT,便于数值分析
pages BIGINT,
price DOUBLE,
isbn STRING,
score DOUBLE,
read_count BIGINT
)
ROW FORMAT DELIMITED
FIELDS TERMINATED BY '\t'
LOCATION '/user/hive/warehouse/dwd.db/dwd_book_detail';
**快速验证**
SELECT * FROM dwd.dwd_book_detail LIMIT 5;
如果长期卡在这个页面并出现Error多半代表这内存不够,(我没有解决)建议切换成本地模式
本地模式:
SET mapreduce.framework.name=local;
SET hive.exec.mode.local.auto=true;
SET hive.exec.mode.local.auto.inputbytes.max=134217728;
SET hive.exec.mode.local.auto.input.files.max=10;
快速验证
SELECT * FROM dwd.dwd_book_detail LIMIT 10;
应该显示数据
八、Python 聚合脚本与 ADS/DWS 表生成
8.1 聚合脚本核心逻辑
脚本路径:/model/agg_dwd_to_app.py
执行:cd /model && python agg_dwd_to_app.py
脚本通过 hdfs dfs -cat 读取 DWD 数据,逐行解析并计算各种聚合指标(出版社、年份、作者、价格区间等),最终将结果写入 HDFS 相应目录。
主要函数说明:
read_dwd_rows():读取 HDFS 上 DWD 目录下的所有文本文件,返回行列表。write_hdfs_table(path, lines):将聚合后的行列表先写入本地临时文件,再上传到 HDFS 指定路径。safe_int()/safe_float():安全类型转换,处理空值等异常情况。
聚合维度举例:
- 出版社维度:统计书籍数量、平均评分、总阅读量、最新出版年份。
- 价格区间:050、50100、100~200、200以上四个档位的数量与平均分。
- 高分书籍:提取评分 ≥ 9.0 的书籍。
- 作者出版年份跨度:计算最早和最晚出版年份差。
完整脚本代码见附录。
运行脚本:
python /model/agg_dwd_to_app.py
8.2 创建指标层外部表
脚本生成以下文件,需在 Hive 中建表映射:
-- DWS 层(汇总宽表)
CREATE DATABASE IF NOT EXISTS dws;
USE dws;
CREATE EXTERNAL TABLE dws_press_stats (
press STRING, book_cnt INT, avg_score DOUBLE, total_reads BIGINT, latest_pub_year INT
) ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t'
LOCATION '/user/hive/warehouse/dws.db/dws_press_stats';
CREATE EXTERNAL TABLE dws_year_stats (
pub_year INT, book_cnt INT, avg_score DOUBLE, total_pages BIGINT, total_reads BIGINT
) ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t'
LOCATION '/user/hive/warehouse/dws.db/dws_year_stats';
CREATE EXTERNAL TABLE dws_author_stats (
author STRING, book_cnt INT, avg_score DOUBLE, total_reads BIGINT
) ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t'
LOCATION '/user/hive/warehouse/dws.db/dws_author_stats';
-- ADS 层(应用指标表)
CREATE DATABASE IF NOT EXISTS ads;
USE ads;
CREATE EXTERNAL TABLE ads_overview (
total_books BIGINT, total_authors BIGINT, total_presses BIGINT,
avg_score DOUBLE, total_reads BIGINT
) ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t'
LOCATION '/user/hive/warehouse/ads.db/ads_overview';
CREATE EXTERNAL TABLE ads_top10_books (
book_name STRING, author STRING, score DOUBLE, read_count BIGINT
) ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t'
LOCATION '/user/hive/warehouse/ads.db/ads_top10_books';
CREATE EXTERNAL TABLE ads_press_ranking (
press STRING, book_cnt INT, avg_score DOUBLE, total_reads BIGINT
) ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t'
LOCATION '/user/hive/warehouse/ads.db/ads_press_ranking';
CREATE EXTERNAL TABLE ads_year_trend (
pub_year INT, book_cnt INT, avg_score DOUBLE, total_reads BIGINT
) ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t'
LOCATION '/user/hive/warehouse/ads.db/ads_year_trend';
CREATE EXTERNAL TABLE ads_hot_authors (
author STRING, book_cnt INT, avg_score DOUBLE, total_reads BIGINT
) ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t'
LOCATION '/user/hive/warehouse/ads.db/ads_hot_authors';
-- 其他 ADS 表(价格区间、页数区间、高分书籍、作者跨度等)类似,参考脚本输出对应创建
8.3 验证查询
SELECT * FROM ads.ads_top10_books;
输出结果
九、总结
本文详细记录了一套完整的离线数仓搭建过程,所有关键配置文件都附带了 行内注释 和 替换指引,即使环境不同也能轻松调整。从集群部署到数据同步,再到多层建模和指标计算,每一步都具备可操作性。
通过本项目,你可以掌握:
- Hadoop 完全分布式集群的核心配置与调优;
- 使用 PostgreSQL 作为 Hive 元数据库的实践;
- DataX 在不同数据源间同步和清洗数据的灵活用法;
- Python 自定义聚合逻辑与 Hive 外部表结合构建分析层。
提示:所有代码和配置文件均通过测试,欢迎复制使用。如遇问题,请检查 IP、密码、路径是否已替换为自己的环境。
附录:完整 Python 聚合脚本
#!/usr/bin/env python
# -*- coding: utf-8 -*-
"""
从 DWD 明细表聚合生成所有 DWS / ADS 表,写入 HDFS。
兼容 Python 2.7 / 3.x
"""
from __future__ import print_function, unicode_literals
import subprocess
import os
import io
from collections import defaultdict, OrderedDict
# ========== 可配置参数 ==========
HDFS_DWD_DIR = "/user/hive/warehouse/dwd.db/dwd_book_detail" # DWD 数据路径
HDFS_BASE = "/user/hive/warehouse" # 仓库根路径
DELIM = "\t" # 输出分隔符
def read_dwd_rows():
"""读取 HDFS 上 DWD 目录下的所有文本文件"""
cmd = "hdfs dfs -cat {0}/*".format(HDFS_DWD_DIR)
proc = subprocess.Popen(cmd, shell=True, stdout=subprocess.PIPE, stderr=subprocess.PIPE)
rows = []
for line in proc.stdout:
line_str = line.decode('utf-8') if isinstance(line, bytes) else line
parts = line_str.strip().split('\t')
if len(parts) >= 10:
rows.append(parts)
proc.wait()
return rows
def write_hdfs_table(file_path, lines):
"""将结果列表写入 HDFS 指定路径"""
local_tmp = "/tmp/agg_tmp.txt"
with io.open(local_tmp, 'w', encoding='utf-8') as f:
for l in lines:
f.write(l + '\n')
subprocess.call(["hdfs", "dfs", "-mkdir", "-p", file_path])
subprocess.call(["hdfs", "dfs", "-put", "-f", local_tmp, file_path + "/data.txt"])
os.remove(local_tmp)
def safe_int(s):
try: return int(float(s))
except: return 0
def safe_float(s):
try: return float(s)
except: return 0.0
def main():
print("正在读取 DWD 数据...")
rows = read_dwd_rows()
print("读取到 {} 行数据".format(len(rows)))
# 聚合容器初始化
press_stats = defaultdict(lambda: {"cnt":0, "score_sum":0.0, "reads_sum":0, "max_year":0})
year_stats = defaultdict(lambda: {"cnt":0, "score_sum":0.0, "pages_sum":0, "reads_sum":0})
author_stats = defaultdict(lambda: {"cnt":0, "score_sum":0.0, "reads_sum":0})
price_range = OrderedDict([
("0~50", [0,50]),
("50~100", [50,100]),
("100~200",[100,200]),
("200以上",[200,99999])
])
pages_range = OrderedDict([
("0~200", [0,200]),
("200~400",[200,400]),
("400以上",[400,99999])
])
price_stats = {k:{"cnt":0, "score_sum":0.0} for k in price_range}
pages_stats = {k:{"cnt":0, "score_sum":0.0} for k in pages_range}
high_score_books = []
author_year_dict = {}
all_ids = set()
all_authors = set()
all_presses = set()
total_score = 0.0
total_reads = 0
book_list = []
for cols in rows:
try:
id = int(cols[0])
book_name = cols[1]
author = cols[2] if cols[2] else "未知"
press = cols[3] if cols[3] else "未知"
pub_year = safe_int(cols[4]) if cols[4] else None
pages = safe_int(cols[5])
price = safe_float(cols[6])
score = safe_float(cols[8])
read_count = safe_int(cols[9])
except:
continue
all_ids.add(id)
all_authors.add(author)
all_presses.add(press)
total_score += score
total_reads += read_count
book_list.append( (score, book_name, author, read_count) )
p = press_stats[press]
p["cnt"] += 1
p["score_sum"] += score
p["reads_sum"] += read_count
if pub_year and pub_year > p["max_year"]:
p["max_year"] = pub_year
if pub_year and pub_year > 0:
y = year_stats[pub_year]
y["cnt"] += 1
y["score_sum"] += score
y["pages_sum"] += pages
y["reads_sum"] += read_count
a = author_stats[author]
a["cnt"] += 1
a["score_sum"] += score
a["reads_sum"] += read_count
for pr_key, (lo, hi) in price_range.items():
if lo <= price < hi:
price_stats[pr_key]["cnt"] += 1
price_stats[pr_key]["score_sum"] += score
break
for pg_key, (lo, hi) in pages_range.items():
if lo <= pages < hi:
pages_stats[pg_key]["cnt"] += 1
pages_stats[pg_key]["score_sum"] += score
break
if score >= 9.0:
high_score_books.append( (book_name, author, score, read_count) )
if pub_year and pub_year > 0:
if author not in author_year_dict:
author_year_dict[author] = {"earliest": pub_year, "latest": pub_year, "cnt": 1}
else:
rec = author_year_dict[author]
if pub_year < rec["earliest"]: rec["earliest"] = pub_year
if pub_year > rec["latest"]: rec["latest"] = pub_year
rec["cnt"] += 1
print("正在生成表文件...")
# dws_press_stats
lines = []
for press, v in press_stats.items():
avg = round(v["score_sum"]/v["cnt"],2) if v["cnt"] else 0.0
lines.append("{}\t{}\t{:.2f}\t{}\t{}".format(press, v["cnt"], avg, v["reads_sum"], v["max_year"]))
write_hdfs_table(HDFS_BASE + "/dws.db/dws_press_stats", lines)
# dws_year_stats
lines = []
for year, v in year_stats.items():
avg = round(v["score_sum"]/v["cnt"],2) if v["cnt"] else 0.0
lines.append("{}\t{}\t{:.2f}\t{}\t{}".format(year, v["cnt"], avg, v["pages_sum"], v["reads_sum"]))
write_hdfs_table(HDFS_BASE + "/dws.db/dws_year_stats", lines)
# dws_author_stats
lines = []
for author, v in author_stats.items():
avg = round(v["score_sum"]/v["cnt"],2) if v["cnt"] else 0.0
lines.append("{}\t{}\t{:.2f}\t{}".format(author, v["cnt"], avg, v["reads_sum"]))
write_hdfs_table(HDFS_BASE + "/dws.db/dws_author_stats", lines)
# ads_overview
total_books = len(all_ids)
total_authors = len(all_authors)
total_presses = len(all_presses)
avg_score_all = round(total_score/total_books,2) if total_books else 0.0
write_hdfs_table(HDFS_BASE + "/ads.db/ads_overview",
["{}\t{}\t{}\t{:.2f}\t{}".format(total_books, total_authors, total_presses, avg_score_all, total_reads)])
# ads_top10_books
book_list.sort(key=lambda x: x[0], reverse=True)
top10 = book_list[:10]
lines = []
for sc, bn, au, rc in top10:
lines.append("{}\t{}\t{:.2f}\t{}".format(bn, au, sc, rc))
write_hdfs_table(HDFS_BASE + "/ads.db/ads_top10_books", lines)
# ads_press_ranking (Top20)
press_rank = sorted(press_stats.items(), key=lambda x: x[1]["cnt"], reverse=True)[:20]
lines = []
for press, v in press_rank:
avg = round(v["score_sum"]/v["cnt"],2) if v["cnt"] else 0.0
lines.append("{}\t{}\t{:.2f}\t{}".format(press, v["cnt"], avg, v["reads_sum"]))
write_hdfs_table(HDFS_BASE + "/ads.db/ads_press_ranking", lines)
# ads_year_trend
years_sorted = sorted(year_stats.items())
lines = []
for year, v in years_sorted:
if 1900 <= year <= 2026:
avg = round(v["score_sum"]/v["cnt"],2) if v["cnt"] else 0.0
lines.append("{}\t{}\t{:.2f}\t{}".format(year, v["cnt"], avg, v["reads_sum"]))
write_hdfs_table(HDFS_BASE + "/ads.db/ads_year_trend", lines)
# ads_hot_authors (Top20)
auth_rank = sorted(author_stats.items(), key=lambda x: x[1]["reads_sum"], reverse=True)[:20]
lines = []
for author, v in auth_rank:
avg = round(v["score_sum"]/v["cnt"],2) if v["cnt"] else 0.0
lines.append("{}\t{}\t{:.2f}\t{}".format(author, v["cnt"], avg, v["reads_sum"]))
write_hdfs_table(HDFS_BASE + "/ads.db/ads_hot_authors", lines)
# ads_price_range_stats
lines = []
for pr_key in price_range:
v = price_stats[pr_key]
avg = round(v["score_sum"]/v["cnt"],2) if v["cnt"] else 0.0
lines.append("{}\t{}\t{:.2f}".format(pr_key, v["cnt"], avg))
write_hdfs_table(HDFS_BASE + "/ads.db/ads_price_range_stats", lines)
# ads_pages_range_stats
lines = []
for pg_key in pages_range:
v = pages_stats[pg_key]
avg = round(v["score_sum"]/v["cnt"],2) if v["cnt"] else 0.0
lines.append("{}\t{}\t{:.2f}".format(pg_key, v["cnt"], avg))
write_hdfs_table(HDFS_BASE + "/ads.db/ads_pages_range_stats", lines)
# ads_high_score_books
high_score_books.sort(key=lambda x: x[2], reverse=True)
lines = []
for bn, au, sc, rc in high_score_books:
lines.append("{}\t{}\t{:.2f}\t{}".format(bn, au, sc, rc))
write_hdfs_table(HDFS_BASE + "/ads.db/ads_high_score_books", lines)
# ads_author_year_span (Top30)
author_spans = []
for author, rec in author_year_dict.items():
span = rec["latest"] - rec["earliest"]
author_spans.append( (author, rec["earliest"], rec["latest"], span, rec["cnt"]) )
author_spans.sort(key=lambda x: x[3], reverse=True)
top30_spans = author_spans[:30]
lines = []
for author, early, late, span, cnt in top30_spans:
lines.append("{}\t{}\t{}\t{}\t{}".format(author, early, late, span, cnt))
write_hdfs_table(HDFS_BASE + "/ads.db/ads_author_year_span", lines)
print("全部表生成完毕!共 12 张表(3 DWS + 9 ADS)")
if __name__ == "__main__":
main()
更多推荐




所有评论(0)