如何让Apache Spark中不同节点访问并共享同一张SQL表?
Spark跨Worker节点Hive表无法访问问题解决
问题场景
- 在仅部署Worker的机器(Worker1)上运行代码创建Hive表
src,本地可正常读取表内容; - 注释表创建逻辑后,在部署了Master和Worker的机器(Worker2)上执行读取代码,触发
TABLE_OR_VIEW_NOT_FOUND错误; - 两台Worker节点均能在Master的8080管理页面被正常识别。
测试代码
""" from pyspark.sql import SparkSession spark = SparkSession \ .builder \ .appName("Python Spark SQL basic example") \ .config("spark.some.config.option", "some-value") \ .getOrCreate() """ from os.path import abspath from pyspark.sql import SparkSession from pyspark.sql import Row # warehouse_location points to the default location for managed databases and tables warehouse_location = abspath('spark-warehouse') #.appName("Python Spark SQL Hive integration example") \ spark = SparkSession \ .builder \ .appName("Python Spark SQL Hive integration") \ .config("spark.sql.warehouse.dir", warehouse_location) \ .enableHiveSupport() \ .getOrCreate() #spark.sql("DROP TABLE src"); #spark.sql("SELECT * FROM src").show() # spark is an existing SparkSession #spark.sql("CREATE TABLE IF NOT EXISTS src (key INT, value STRING) USING hive") #spark.sql("LOAD DATA LOCAL INPATH 'examples/src/main/resources/kv1.txt' INTO TABLE src") #spark.sql("INSERT INTO src VALUES (10,\"HELLO\")") #returnSql = spark.sql("SHOW TABLES").show() #try: #returnSql = spark.sql("DESCRIBE src1") #except: # print("FAIL") #returnSql = spark.sql("SHOW TABLES") returnSql = spark.sql("DESCRIBE src") #spark.collect()[Row number][Column number][0][0] print("----------------------------------------") returnSql.show() print("----------------------------------------") print("****************************************") print(returnSql[0][1]) print(returnSql[1][0]) print(returnSql[2]) #print(returnSql[3]) #print(returnSql[4]) #print(returnSql[5]) print("****************************************") print(returnSql.collect()[0][0]) print("****************************************") #returnSql = spark.sql("SHOW TABLES").show() #spark.sql("SELECT COUNT(*) FROM src").show() import sys # In some cases, you may need to import the sys module sys.exit() # Queries are expressed in HiveQL spark.sql("SELECT * FROM src").show() # +---+-------+ # |key| value| # +---+-------+ # |238|val_238| # | 86| val_86| # |311|val_311| # ... # Aggregation queries are also supported. spark.sql("SELECT COUNT(*) FROM src").show() # +--------+ # |count(1)| # +--------+ # | 500 | # +--------+ # The results of SQL queries are themselves DataFrames and support all normal functions. sqlDF = spark.sql("SELECT key, value FROM src WHERE key < 10 ORDER BY key") spark.sql("SELECT * FROM records r JOIN src s ON r.key = s.key").show() # The items in DataFrames are of type Row, which allows you to access each column by ordinal. #stringsDS = sqlDF.rdd.map(lambda row: "Key: %d, Value: %s" % (row.key, row.value)) #for record in stringsDS.collect(): # print(record) # Key: 0, Value: val_0 # Key: 0, Value: val_0 # Key: 0, Value: val_0 # ... # You can also use DataFrames to create temporary views within a SparkSession. #Record = Row("key", "value") #recordsDF = spark.createDataFrame([Record(i, "val_" + str(i)) for i in range(1, 101)]) #recordsDF.createOrReplaceTempView("records") # Queries can then join DataFrame data with data stored in Hive. #spark.sql("SELECT * FROM records r JOIN src s ON r.key = s.key").show() # +---+------+---+------+ # |key| value|key| value| # +---+------+---+------+ # | 2| val_2| 2| val_2| # | 4| val_4| 4| val_4| # | 5| val_5| 5| val_5| # ...
问题根源
- 本地仓库路径孤立:代码中
warehouse_location = abspath('spark-warehouse')指向节点本地磁盘目录,Worker1创建的表数据仅存储在自身本地目录,Worker2本地无对应表文件; - 元数据未集中管理:未配置外部共享元数据库时,Spark默认使用本地Derby数据库存储Hive元数据,每个节点的Derby实例独立,Worker2的元数据库中无
src表的记录。
解决方案
1. 配置共享Warehouse存储
将spark.sql.warehouse.dir指定为所有节点可访问的共享存储路径,比如HDFS或NFS挂载目录:
# 替换为实际共享路径,示例为HDFS路径 warehouse_location = "hdfs://master:9000/spark-warehouse"
确保所有Worker节点对该路径有读写权限。
2. 配置集中式Hive元数据库
部署共享元数据库(如MySQL),并在所有Spark节点的spark-defaults.conf中添加以下配置(替换为实际环境信息):
spark.sql.hive.metastore.version = 3.1.2 spark.sql.hive.metastore.jars = builtin spark.sql.hive.metastore.sharedPrefixes = com.mysql.cj.jdbc spark.hadoop.javax.jdo.option.ConnectionURL = jdbc:mysql://master:3306/hive_metastore?createDatabaseIfNotExist=true&useSSL=false&serverTimezone=UTC spark.hadoop.javax.jdo.option.ConnectionDriverName = com.mysql.cj.jdbc.Driver spark.hadoop.javax.jdo.option.ConnectionUserName = hive_user spark.hadoop.javax.jdo.option.ConnectionPassword = hive_pass
保证所有节点的配置一致,且能连接到元数据库。
3. 验证配置
在任意Worker节点创建表,然后在其他节点执行读取操作,确认表可正常访问;同时检查元数据库中的表记录,确保元数据已同步。
内容的提问来源于stack exchange,提问作者Rick C. Ferreira
相关产品推荐
相关产品推荐

