You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何让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|
# ...

问题根源

  1. 本地仓库路径孤立:代码中warehouse_location = abspath('spark-warehouse')指向节点本地磁盘目录,Worker1创建的表数据仅存储在自身本地目录,Worker2本地无对应表文件;
  2. 元数据未集中管理:未配置外部共享元数据库时,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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 21:24:54