目录


一、架构与节点规划

本项目模拟图书数据仓库,从业务库 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

根据节点添加masterIPADDR 为 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
master
node1/node2
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()

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐