SparkCore、SparkSQL 数据读取方式及 Hive 数据来源详解

在大数据处理领域,Apache Spark 和 Apache Hive 是两款至关重要的工具。Spark 凭借其高效的内存计算能力,成为数据处理的热门选择,而 Hive 则在数据仓库构建和 SQL 分析方面发挥着重要作用。本文将详细介绍 SparkCore 和 SparkSQL 读取数据的多种方式,以及 Hive 中数据的主要来源,帮助读者更好地理解和运用这些工具进行数据处理。

一、SparkCore 读取数据的方式

SparkCore 作为 Spark 的核心模块,提供了丰富的数据读取接口,支持从多种数据源读取数据,满足不同场景下的数据处理需求。以下是 SparkCore 常见的几种数据读取方式:

1. 从本地文件系统读取数据

本地文件系统是最基础的数据存储位置之一,SparkCore 可以轻松读取本地文件中的数据。这种方式适用于数据量较小,且数据存储在 Spark 运行节点本地的场景。

方法:textFile、wholeTextFile、newAPIHadoopRDD等

读取外部存储系统的数据转换为RDD

例如,我们有一个存储在本地 ./data/test.txt 文件,文件内容为每行一个字符串。使用 SparkCore 读取该文件的代码如下:

	os.environ['JAVA_HOME'] = 'D:/software/JDK/jdk1.8'
    # 配置Hadoop的路径,就是前面解压的那个路径
    os.environ['HADOOP_HOME'] = 'D:/software/hadoop-3.3.1'
    # 配置base环境Python解析器的路径
    os.environ['PYSPARK_PYTHON'] = 'D:/software/Miniconda3/python.exe'
    os.environ['PYSPARK_DRIVER_PYTHON'] = 'D:/software/Miniconda3/python.exe'
    
    conf = SparkConf().setMaster("local[2]").setAppName("单词统计")
    sc = SparkContext(conf=conf)
    print(sc)
	
    textFileRdd = sc.textFile("./datas/test.txt")

在上述代码中,textFile 方法用于读取文本文件。通过这种方式,SparkCore 会将文件内容拆分成多个分区,生成弹性分布式数据集(RDD),以便进行后续的并行计算。

2. 并行化一个已存在的集合

方法:parallelize 并行的意思

将一个集合转换为RDD


# 方式一:将一个已存在的集合转换为RDD
# 创建一个列表:会在Driver内存中构建
data = [1,2,3,4,5,6,7,8,9,10]
# 将列表转换为RDD:将在多个Executor内存中实现分布
式存储, numSlices用于指定分区数,所谓的分区就是分为几份,每一份放在一台电脑上
list_rdd = sc.parallelize(data,numSlices=2)
# 打印这个RDD的内容
list_rdd.foreach(lambda x: print(x))

二、SparkSQL 读取数据的方式

SparkSQL 是 Spark 中用于处理结构化数据的模块,它提供了 SQL 接口和 DataFrame/Dataset API,支持从多种数据源读取结构化数据,简化了数据的查询和分析过程。以下是 SparkSQL 常见的几种数据读取方式:

1. 从文本文件读取数据(TEXT、CSV、JSON 等格式)

SparkSQL 支持读取多种格式的文本文件,如text 、CSV、JSON等,并且能够自动推断数据的 schema(结构),大大简化了数据读取的过程。

(1)方式一:给定读取数据源的类型和地址
spark.read.format("json").load(path)
spark.read.format("csv").load(path)
spark.read.format("parquet").load(path)

通过 csv 方法读取 CSV 文件,并通过 option 方法设置相关参数。读取后的数据以 DataFrame 的形式存在,可以方便地进行 SQL 查询和数据分析。

(2)方式二:直接调用对应数据源类型的方法
spark.read.json(path)
spark.read.csv(path)
spark.read.parquet(path)

2.读取数据库数据

与 SparkCore 类似,SparkSQL 也可以通过 JDBC 连接关系型数据库,读取数据库中的表数据,并将其转换为 DataFrame,方便进行 SQL 分析。

以从 MySQL 数据库读取 emp 表数据为例,代码如下:

empDf= spark.read.format("jdbc").option("url", "jdbc:mysql://localhost:3306/spark_demo") \
            .option("dbtable", "emp") \
            .option("user", "root") \
            .option("password", "123456") \
            .load()
        empDf.show()

SparkSQL 提供的 read.jdbc 方法简化了从数据库读取数据的过程,并且可以直接对生成的 DataFrame 进行各种数据操作和分析,与 SQL 语法高度兼容。

3. 从 Hive 表读取数据

SparkSQL 可以与 Hive 集成,直接读取 Hive 表中的数据,无需额外的数据迁移操作。这种方式适用于已经构建了 Hive 数据仓库,需要利用 Spark 的计算能力进行更复杂数据分析的场景。

要实现 SparkSQL 读取 Hive 表数据,需要在 Spark 的配置文件(如 hive-site.xml)中配置 Hive 的相关信息,确保 Spark 能够连接到 Hive metastore。以下是读取 Hive 表数据的示例代码:

spark = SparkSession \
        .builder \
        .appName("HiveAPP") \
        .master("local[2]") \
        .config("spark.sql.warehouse.dir", 'hdfs://bigdata01:9820/user/hive/warehouse') \
        .config('hive.metastore.uris', 'thrift://bigdata01:9083') \
        .config("spark.sql.shuffle.partitions", 2) \
        .enableHiveSupport()\
        .getOrCreate()

代码实战:

from pyspark.sql import SparkSession

if __name__ == '__main__':
    spark = SparkSession \
        .builder \
        .appName("sparksql操作hive") \
        .master("local[2]") \
        .config("spark.sql.warehouse.dir", 'hdfs://caijing:9820/user/hive/warehouse') \
        .config('hive.metastore.uris', 'thrift://caijing:9083') \
        .config("spark.sql.shuffle.partitions", 2) \
        .enableHiveSupport() \
        .getOrCreate()

    spark.sql("""
        select * from caijing01.t_user
    """).show()

    spark.stop()

代码还可以这样写:

方式二:加载Hive表的数据变成DF,可以调用DSL或者SQL的方式来实现计算

# 读取Hive表构建DataFrame

hiveData = spark.read.table(“yhdb.student”)

hiveData.printSchema()

hiveData.show()

# 读取hive表中的数据
    os.environ['HADOOP_USER_NAME'] = 'root'
	spark2 = SparkSession \
		.builder \
		.appName("HiveAPP") \
		.master("local[2]") \
		.config("spark.sql.warehouse.dir", 'hdfs://192.168.233.128:9820/user/hive/warehouse') \
		.config('hive.metastore.uris', 'thrift://192.168.233.128:9083') \
		.config("spark.sql.shuffle.partitions", 2) \
		.enableHiveSupport() \
		.getOrCreate()

	#spark2.sql("show databases").show()
	#spark2.sql("show  tables").show()

	#spark2.sql("select * from yhdb.t_user").show()

	spark2.read.table("t_user2").show()

不要在一个python 文件中,创建两个不同的sparkSession对象,否则对于sparksql获取hive的元数据,有影响。另外,记得添加一个权限校验的语句:

# 防止在本地操作hdfs的时候,出现权限问题
os.environ['HADOOP_USER_NAME'] = 'root'

三、Hive 中数据的来源

Hive 作为构建在 Hadoop 之上的数据仓库工具,其数据来源非常广泛,主要包括以下几种途径:

1. 从本地文件系统或 HDFS 加载数据

本地文件系统和 HDFS 是 Hive 最常见的数据来源之一。用户可以将本地文件或 HDFS 上的文件加载到 Hive 表中,以便进行后续的数据分析。

(1)从本地文件系统加载数据

先有数据,根据数据的格式,和字段数量以及类型,创建一个表:

create table t_user(
id int,
name string
)
row format delimited
fields terminated by ','
lines terminated by '\n'
stored as textfile;

加载本地数据:

load data local inpath "/home/hivedata/user.txt" into table t_user;

查看数据是否加载成功:
select * from t_user limit 10;

在上述命令中,ROW FORMAT DELIMITED FIELDS TERMINATED BY ',' 指定了数据的分隔符为逗号。执行 LOAD DATA 命令后,Hive 会将本地文件中的数据复制到 Hive 表对应的存储目录下(通常在 HDFS 上)。

(2)从 HDFS 加载数据

如果数据已经存储在 HDFS 上,可以使用 LOAD DATA INPATH 命令将其加载到 Hive 表中。与从本地加载数据不同,从 HDFS 加载数据时,Hive 会将数据移动(而不是复制)到表的存储目录下(除非使用 COPY 选项)。示例代码如下:

load data inpath '/home/user.txt' into table t_user;

就比之前少了一个 local 关键字。

查看hdfs上的user.txt 发现不见了,去哪里了,被移动走了!

思考一下:为什么是移动,而不是复制?

因为hdfs上的数据,默认都以为比较大,所以如果相同的数据占2份,非常的消耗空间。

查看数据,发现有两份,想覆盖怎么办?

load data local inpath '/home/hivedata/user.txt' overwrite into table t_user;
(3) 将数据直接放入表对应的文件夹下

再思考一个问题:既然hive中的数据是在hdfs上的,我们也可以手动的上传数据,能上传至/home,为何不能上传至:/user/hive/warehouse/yhdb.db/t_user

[root@bigdata01 hivedata]# cp user.txt user2.txt 
[root@bigdata01 hivedata]# hdfs dfs -put /home/hivedata/user2.txt /user/hive/warehouse/yhdb.db/t_user

hive中的数据,不要load 也可以被正常使用。

(4)从其他表中加载数据

语法格式:

insert into table tableName2 select [.....] from tableName1;

扩展内容:向多张表中插入数据的语法
    from tableName1
    insert into tableName2 select * where 条件
    insert into tableName3 select * where 条件

实战:
insert into table t_user2 select * from t_user;
这个sql的前提条件是:必须先创建一个t_user2
快速创建一个同样的表,只要表结构:
create table t_user2 like t_user;
创建完之后再运行
insert into table t_user2 select * from t_user;


创建t_user3和 4
create table t_user5 like t_user;
create table t_user4 like t_user;

from t_user2
insert into t_user4 select *
insert into t_user5 select id,name;
(5) 克隆表数据
- create table if not exists tableName2 as select [....] from tableName1;
- create table if not exists tableName2 like tableName1 location 'tableName1的存储目录的路径'     # 新表不会产生自己的表目录,因为用的是别的表的路径
​
扩展内容:只复制表结构
create table if not exists tableName2 like tableName1;

实战:
create table t_user6 as select * from t_user2;
create table t_user7 like t_user2 location '/user/hive/warehouse/hive02.db/t_user2';

2. 通过 ETL 工具导入数据

在实际的大数据处理流程中,通常需要对原始数据进行抽取(Extract)、转换(Transform)和加载(Load),即 ETL 过程。

常见的ETL 工具 及其抽取数据到hive 代码示例如下:

(1)Flume 抽取数据到 hive
a1.sources = r1
a1.channels = c1
a1.sources.r1.type = exec
a1.sources.r1.command = tail -f /home/hivedata/user.txt
a1.sources.r1.channels = c1

a1.channels.c1.type = memory
a1.channels.c1.capacity = 10000
a1.channels.c1.transactionCapacity = 10000
a1.channels.c1.byteCapacityBufferPercentage = 20
a1.channels.c1.byteCapacity = 800000


a1.sinks = k1
a1.sinks.k1.type = hive
a1.sinks.k1.channel = c1
a1.sinks.k1.hive.metastore = thrift://bigdata01:9083
a1.sinks.k1.hive.database = mydb03
a1.sinks.k1.hive.table = flume_user

a1.sinks.k1.serializer = DELIMITED
a1.sinks.k1.serializer.delimiter = ","
a1.sinks.k1.serializer.serdeSeparator = ','
a1.sinks.k1.serializer.fieldnames =id,name
(2)sqoop :MySQL数据导入到Hive
sqoop import --connect jdbc:mysql://bigdata01:3306/sqoop \
--driver com.mysql.cj.jdbc.Driver \
--username root \
--password 123456 \
--table emp \
--hive-import \
--hive-overwrite \
--hive-table emp \
--hive-database yhdb \
-m 1
(3)DataX :MySQL数据导入到Hive

首先在hive中创建一个ods_01_base_area

create external table if not exists ods_01_base_area (
  id int COMMENT 'id标识',
  area_code string COMMENT '省份编码',
  province_name string COMMENT '省份名称',
  iso string COMMENT 'ISO编码'
)row format delimited fields terminated by ','
stored as TextFile
location '/data/nshop/ods/ods_01_base_area/';

从mysql导入到hive中(其实就是导入到hdfs)

编写对应的Job的json文件:

read 方 是mysqlreader

write 方 是 hive (没有找到,找到了hdfswriter)

在 job文件夹,创建一个 mysql2hive01.json

{
    "job": {
        "setting": {
            "speed": {
                 "channel": 3
            },
            "errorLimit": {
                "record": 0,
                "percentage": 0.02
            }
        },
        "content": [
            {
                "reader": {
                    "name": "mysqlreader",
                    "parameter": {
                        "username": "root",
                        "password": "123456",
                        "column": [
                            "id",
                            "area_code",
                            "province_name",
                            "iso"
                        ],
                        "splitPk": "id",
                        "connection": [
                            {
                                "table": [
                                    "base_area"
                                ],
                                "jdbcUrl": [
     "jdbc:mysql://bigdata01:3306/sqoop"
                                ]
                            }
                        ]
                    }
                },
               "writer": {
                    "name": "hdfswriter",
                    "parameter": {
                        "defaultFS": "hdfs://bigdata01:9820",
                        "fileType": "text",
                        "path": "/data/nshop/ods/ods_01_base_area/",
                        "fileName": "base_area_txt",
                        "column": [
                            {
                                "name": "id",
                                "type": "int"
                            },
                            {
                                "name": "area_code",
                                "type": "string"
                            },
                            {
                                "name": "province_name",
                                "type": "string"
                            },
                            {
                                "name": "iso",
                                "type": "string"
                            }
                        ],
                        "writeMode": "append",
                        "fieldDelimiter": ","
                    }
                }
            }
        ]
    }
}

更多推荐