基於Hologres提供的Spark Connector,EMR Serverless Spark可以在開發時添加對應的配置來串連Hologres。本文為您介紹在EMR Serverless Spark環境中實現Hologres的資料讀取和寫入操作。
使用限制
僅V1.3及以上版本的Hologres執行個體支援Spark Connector。您可以在Hologres管理主控台的執行個體詳情頁查看當前執行個體版本。若您的執行個體是V1.3以下版本,請使用執行個體升級或通過搜尋(DingTalk群號:32314975)加入即時數倉Hologres交流群申請升級執行個體。
訪問方式
訪問 Hologres 有兩種方式,您可以根據實際需求選擇:
訪問方式 | 說明 | 適用情境 | 參考文檔 |
方式一:任務/會話層級配置 | 需要在每個任務或會話中單獨配置 Hologres 的串連資訊(包括 JDBC URL、使用者名稱、密碼等)。 |
| 本文檔 |
方式二:通過資料目錄統一配置(推薦) | 通過 EMR Serverless Spark 的資料目錄功能添加 Hologres 資料目錄。添加後,該工作空間下提交的所有作業和建立的會話預設可以訪問該資料目錄下有許可權的資料,無需在每個任務中重複配置。 說明 僅支援使用以下引擎版本:esr-4.9.0 及以上版本。 |
|
如果您的工作空間需要長期、頻繁地訪問 Hologres 資料,推薦使用方式二(資料目錄),可以減少重複配置,提升開發效率。
操作流程
步驟一:擷取 hologres-connector-spark JAR並上傳至OSS
esr-4.8.0及以上版本的引擎已內建Hologres連接器,無需執行此步驟。
Spark讀寫Hologres需要引用的連接器JAR包,您可以通過Maven中央倉庫進行下載。本文提供1.5.6版本,您可以通過單擊附件hologres-connector-spark-3.x-1.5.6-jar-with-dependencies.jar下載。
將下載的 hologres-connector-spark JAR上傳至阿里雲OSS中,上傳操作可以參見簡單上傳。
步驟二:添加網路連接
擷取網路資訊。
您可以在即時數倉Hologres頁面,進入目標Hologres執行個體的執行個體詳情頁面,以擷取該執行個體的專用網路和交換器資訊。
新增網路連接。
Serverless Spark需要能夠打通與Hologres叢集之間的網路才可以正常訪問Hologres服務。有關更多網路連接資訊,請參見EMR Serverless Spark與其他VPC間網路互連。
步驟三:在Hologres中建立庫表
串連Hologres執行個體,詳情請參見串連執行個體。
在SQL編輯器頁簽,新增的臨時Query查詢中輸入以下SQL語句,並執行。
-- 建立資料庫 CREATE DATABASE testdb; -- 建立表 CREATE TABLE "public"."test" ( "id" text NULL, "name" text NULL); -- 插入資料 INSERT INTO public.test VALUES ('1001','jack'),('1002','tony'),('1003','mike'); -- 查詢資料 SELECT * FROM public.test
步驟四:通過Serverless Spark讀寫Hologres
樣本1:SQL會話
以SQL會話為例,讀寫Hologres。
建立SQL會話,詳情請參見管理SQL會話。
建立會話時,在網路連接中選擇上一步建立好的網路連接,並在Spark 配置中添加以下參數來載入hologres-connector-spark。
# 添加hologres-connector jar(僅適用於esr-4.8.0以下版本,esr-4.8.0及以上版本無需配置) spark.emr.serverless.user.defined.jars oss://<bucket>/hologres-connector-spark-3.x-<version>.jar # 配置holo catalog spark.sql.catalog.hologres_external_test_db com.alibaba.hologres.spark3.HoloTableCatalog spark.sql.catalog.hologres_external_test_db.username *** spark.sql.catalog.hologres_external_test_db.password *** spark.sql.catalog.hologres_external_test_db.jdbcurl jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/testdb參數詳細說明如下:
參數
樣本
說明
spark.emr.serverless.user.defined.jarsoss://<bucket>/hologres-connector-spark-3.x-<version>.jar指定使用者自訂的 JAR 包路徑。
spark.sql.catalog.hologres_external_test_dbcom.alibaba.hologres.spark3.HoloTableCatalogSpark 3.x 中用於配置 Hologres 資料來源作為外部 Catalog,固定值。
spark.sql.catalog.hologres_external_test_db.usernameLTAI******阿里雲帳號的AccessKey ID。推薦採用密文方式管理敏感資訊,詳細資料請參見通過密文管理敏感資訊。
spark.sql.catalog.hologres_external_test_db.passwordmXYV******阿里雲帳號的AccessKey Secret。推薦採用密文方式管理敏感資訊,詳細資料請參見通過密文管理敏感資訊。
spark.sql.catalog.hologres_external_test_db.jdbcurljdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/test_dbHologres 執行個體 JDBC 串連 URL 。
其中參數名中
hologres_external_test_db可自訂。在資料開發頁面,建立一個SparkSQL類型的任務,然後在右上方選擇建立好的SQL會話。
更多操作,請參見SparkSQL開發。
拷貝如下代碼到新增的SparkSQL頁簽中,然後單擊运行。
-- 進入testdb database USE hologres_external_test_db; -- 寫入資料 INSERT INTO `public`.test VALUES ('1004','tom'); -- 查詢資料 SELECT * FROM `public`.test;執行查詢後返回 4 條記錄:
1001, jack、1002, tony、1003, mike、1004, tom。
樣本2:流任務
以流任務PySpark為例,從Kafka讀資料,寫入Hologres。
確保Kafka與Hologres之間的網路連接暢通,建議將Kafka與Hologres部署在同一個VPC及同一交換器中。
程式碼範例。根據實際情況替換Kafka資訊和Hologres表。
from pyspark.sql import SparkSession from pyspark.sql.functions import col # 配置你的 Kafka 資訊 servers = "alikafka-serverless-cn-xxxxx-vpc.alikafka.aliyuncs.com:9092" # 替換為你的 Kafka bootstrap servers topic = "topic-name" # 替換為你的 Kafka topic # 建立 SparkSession spark = SparkSession.builder \ .appName("test read kafka") \ .getOrCreate() # 讀取 Kafka 流 df = spark \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", servers) \ .option("subscribe", topic) \ .load() # 定義寫入 Hologres 的函數(每個 micro-batch 調用一次) def write_to_hologres(batch_df, batch_id): print(f"Writing batch {batch_id} to Hologres...") batch_df.write \ .format("hologres") \ .mode("append") \ .insertInto("hologres_external_test_db.public.test") # 替換為你的 Hologres table # 轉換 key 和 value 為字串並輸出到控制台 query = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") \ .writeStream \ .foreachBatch(write_to_hologres) \ .outputMode("append") \ .trigger(processingTime='30 seconds') \ .start() # 等待流式查詢結束(在 notebook 中會阻塞 cell 執行) query.awaitTermination()上傳檔案。
在文件管理頁面,單擊上传文件。
在上传文件對話方塊中,單擊待上傳檔案地區選取項目上一步的Python代碼檔案,或直接拖拽Python代碼檔案到待上傳檔案地區。
建立流任務並運行。
在数据开发頁面,單擊
(建立)表徵圖。在彈出的對話方塊中,輸入名称,根據實際需求在流任务中選擇PySpark類型,然後單擊确定。
在建立的開發頁簽中,配置以下資訊,其餘參數無需配置,然後單擊发布。
參數
說明
主 Python 资源
選擇前一個步驟中上傳的Python檔案。
引擎版本
選擇合適的Spark版本。本文樣本是
esr-4.6.0。网络连接
選擇步驟二中建立的網路。
Spark 配置
# 添加hologres-connector jar(僅適用於esr-4.8.0以下版本,esr-4.8.0及以上版本無需配置) spark.emr.serverless.user.defined.jars oss://<bucket>/test_script/hologres-connector-spark-3.x-1.5.6-jar-with-dependencies.jar # 配置holo catalog spark.sql.catalog.hologres_external_test_db com.alibaba.hologres.spark3.HoloTableCatalog spark.sql.catalog.hologres_external_test_db.username *** spark.sql.catalog.hologres_external_test_db.password *** spark.sql.catalog.hologres_external_test_db.jdbcurl jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/testdb參數的詳細說明請參見樣本1:SQL會話。
發布後,單擊前往运维,在跳轉頁面,單擊启动。
驗證結果。
Kafka發送訊息。
在發送方式地區選取項目控制台,在訊息 Key中輸入
1005,在訊息內容中輸入zy。SparkSQL查詢。查詢結果返回5條記錄,包含id和name兩列,資料分別為:1001/jack、1002/tony、1003/mike、1005/zy、1004/tom。
樣本3:Notebook會話
建立Notebook會話,詳情請參見管理SQL會話。
建立會話時,在網路連接中選擇上一步建立好的網路連接,並在Spark 配置中添加以下參數來載入hologres-connector-spark。
# 添加hologres-connector jar(僅適用於esr-4.8.0以下版本,esr-4.8.0及以上版本無需配置) spark.emr.serverless.user.defined.jars oss://<bucket>/hologres-connector-spark-3.x-<version>.jar在資料開發頁面,建立一個Notebook類型的任務,然後在右上方選擇建立好的Notebook會話。
拷貝如下代碼到新增的Notebook頁簽中,然後單擊
。import pandas as pd from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, IntegerType, LongType # 1. 準備 Pandas DataFrame(修複 hsap 為字串) pdf = pd.DataFrame({ "id": ["1006"], # 改為整數,匹配 Hologres 的 BIGINT/INT "name": ["sl"] # 字串 }) # 2. 轉換為 PySpark DataFrame(可選:顯式定義 schema 以確保類型正確) schema = StructType([ StructField("id", StringType(), True), StructField("name", StringType(), True) ]) df = spark.createDataFrame(pdf, schema=schema) # 寫入hologres df.write \ .format("hologres") \ .option("username", "LTAI******") \ .option("password", "mXYV******") \ .option("jdbcurl", "jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/test") \ .option("table", "test") \ .mode("append") \ .save() # 讀取資料 readDf = spark.read\ .format("hologres") \ .option("username", "LTAI******") \ .option("password", "mXYV******") \ .option("jdbcurl", "jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/test") \ .option("table", "test") \ .load() readDf.select("id", "name").show(10)其中,涉及參數說明如下:
參數
樣本
說明
spark.sql.catalog.hologres_external_test_db.usernameLTAI******阿里雲帳號的AccessKey ID。推薦採用密文方式管理敏感資訊,詳細資料請參見通過密文管理敏感資訊。
spark.sql.catalog.hologres_external_test_db.passwordmXYV******阿里雲帳號的AccessKey Secret。推薦採用密文方式管理敏感資訊,詳細資料請參見通過密文管理敏感資訊。
spark.sql.catalog.hologres_external_test_db.jdbcurljdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/test_dbHologres 執行個體 JDBC 串連 URL 。
驗證結果。執行上述代碼後,返回如下查詢結果。
+----+-----+ | id|name | +----+-----+ |1001|jack | |1002|tony | |1003|mike | |1005| zy | |1004| tom | |1006| sl | +----+-----+
Hologres Catalog常用命令
在 Serverless Spark 中,通過 Hologres Catalog 將 Hologres 資料庫(Database)以 外部 Catalog 形式接入 Spark SQL。每個 Catalog 嚴格綁定一個 Hologres Database,且不可跨庫訪問(即無法通過同一 Catalog 訪問多個 Hologres DB)。Catalog 內部的邏輯組織與 Hologres 保持一致:
Spark 概念 | 對應 Hologres 概念 | 說明 |
Catalog | Database | 如 |
Namespace | Schema | 如 |
Table | Table | 必須顯式指定 |
載入Hologres Catalog
Spark中的Hologres Catalog完全對應一個Hologres的Database,使用過程中無法更改。
USE hologres_external_test_db;查詢所有Namespace
Spark中的Namespace,對應Hologres中的Schema,預設為public,使用過程中可以使用USE指令調整預設的Schema。
-- 查看Hologres Catalog中的所有Namespace, 即Hologres中所有的Schema。
SHOW NAMESPACES;查詢Namespace下的表
查詢所有表。
SHOW TABLES;查詢指定Namespace下的表。
USE test_schema; SHOW TABLES; -- 或者使用 SHOW TABLES IN test_schema;
相關文檔
有關Spark讀寫Hologres的更多資訊,請參見Spark讀寫Hologres。