Untuk menjalankan kueri Spark SQL secara interaktif, Anda dapat menentukan kelompok sumber daya Spark Interactive. Kelompok sumber daya tersebut akan melakukan penskalaan otomatis dalam rentang yang ditentukan guna memenuhi kebutuhan analisis interaktif Anda sekaligus mengurangi biaya. Topik ini menjelaskan cara melakukan analisis interaktif dengan Spark SQL menggunakan Konsol, Hive JDBC, PyHive, Beeline, DBeaver, dan tool client lainnya.
Prasyarat
Kluster AnalyticDB for MySQL Edisi Perusahaan, Edisi Dasar, atau Edisi Data Lakehouse telah dibuat.
Bucket Object Storage Service (OSS) telah dibuat di wilayah yang sama dengan kluster AnalyticDB for MySQL.
Akun database telah dibuat untuk kluster AnalyticDB for MySQL.
Jika Anda menggunakan akun Alibaba Cloud, Anda hanya perlu membuat akun istimewa.
Jika Anda menggunakan pengguna Resource Access Management (RAM), Anda harus membuat akun istimewa dan akun standar serta mengaitkan akun standar tersebut dengan pengguna RAM.
-
Lingkungan pengembangan Java 8 dan Python 3.9 harus diinstal untuk menjalankan client seperti aplikasi Java, aplikasi Python, dan Beeline.
-
Alamat IP client Anda harus berada dalam AnalyticDB for MySQL kluster daftar putih.
Catatan penggunaan
-
Jika kelompok sumber daya Spark Interactive dihentikan, kluster akan memulai ulang saat Anda menjalankan kueri Spark SQL pertama. Kueri pertama mungkin masuk antrian selama startup.
-
Spark tidak dapat membaca dari atau menulis ke database INFORMATION_SCHEMA dan MYSQL. Jangan gunakan database tersebut sebagai database koneksi awal.
-
Pastikan akun database yang digunakan untuk mengirim pekerjaan Spark SQL memiliki akses ke database target. Jika tidak, kueri akan gagal.
Persiapan
-
Dapatkan titik akhir kelompok sumber daya Spark Interactive.
Masuk ke Konsol AnalyticDB for MySQL. Di pojok kiri atas konsol, pilih wilayah. Di panel navigasi sebelah kiri, klik Clusters. Temukan kluster yang ingin Anda kelola lalu klik ID kluster tersebut.
-
Di panel navigasi sebelah kiri, pilih , lalu klik tab Resource Groups.
-
Temukan kelompok sumber daya tersebut lalu klik Details pada kolom Actions untuk melihat titik akhir internal dan titik akhir publik. Anda dapat mengklik ikon
di samping titik akhir untuk menyalinnya, atau mengklik ikon
di dalam tanda kurung nomor Port untuk menyalin string koneksi JDBC.Dalam kasus berikut, Anda harus mengklik Apply for Endpoint di samping Public Address untuk mengajukan titik akhir publik secara manual.
-
Tool client yang digunakan untuk mengirim pekerjaan Spark SQL dideploy di mesin lokal Anda atau server eksternal.
-
Tool client yang digunakan untuk mengirim pekerjaan Spark SQL dideploy di instans ECS, dan instans ECS tersebut serta kluster AnalyticDB for MySQL tidak berada dalam VPC yang sama.
Informasi koneksi juga mencakup bidang-bidang seperti port publik dan VPC (default:
10000), ID VPC, ID VSwitch, kelas driver (org.apache.hive.jdbc.HiveDriver), dan URL unduhan driver. -
Analisis interaktif
Konsol
Jika Anda menggunakan HiveMetastore yang dikelola sendiri, buat database bernama default di AnalyticDB for MySQL dan pilih database tersebut sebagai database saat menjalankan pekerjaan Spark SQL di konsol.
Masuk ke Konsol AnalyticDB for MySQL. Di pojok kiri atas konsol, pilih wilayah. Di panel navigasi sebelah kiri, klik Clusters. Temukan kluster yang ingin Anda kelola lalu klik ID kluster tersebut.
-
Di panel navigasi sebelah kiri, pilih .
-
Pilih engine Spark dan kelompok sumber daya Spark Interactive yang telah dibuat, lalu jalankan pernyataan Spark SQL berikut:
SHOW DATABASES;
SDK
Saat Anda menjalankan pernyataan Spark SQL menggunakan SDK, hasil kueri akan ditulis sebagai file ke bucket OSS yang ditentukan. Anda kemudian dapat mengkueri data tersebut di konsol OSS atau mengunduh file hasil ke komputer Anda. Contoh berikut menunjukkan cara memanggil SDK dalam Python.
-
Jalankan perintah berikut untuk menginstal SDK.
pip install alibabacloud-adb20211201 -
Jalankan perintah berikut untuk menginstal dependensi.
pip install oss2 pip install loguru -
Hubungkan ke kluster dan jalankan pernyataan Spark SQL.
# coding: utf-8 import csv import json import time from io import StringIO import oss2 from alibabacloud_adb20211201.client import Client from alibabacloud_adb20211201.models import ExecuteSparkWarehouseBatchSQLRequest, ExecuteSparkWarehouseBatchSQLResponse, \ GetSparkWarehouseBatchSQLRequest, GetSparkWarehouseBatchSQLResponse, \ ListSparkWarehouseBatchSQLRequest, CancelSparkWarehouseBatchSQLRequest, ListSparkWarehouseBatchSQLResponse from alibabacloud_tea_openapi.models import Config from loguru import logger def build_sql_config(oss_location, spark_sql_runtime_config: dict = None, file_format = "CSV", output_partitions = 1, sep = "|"): """ Membuat konfigurasi untuk eksekusi SQL AnalyticDB for MySQL. :param oss_location: Jalur OSS tempat menyimpan hasil eksekusi SQL. :param spark_sql_runtime_config: Properti konfigurasi native Spark SQL. :param file_format: Format file hasil eksekusi SQL. Nilai default: CSV. :param output_partitions: Jumlah partisi untuk hasil eksekusi SQL. Jika Anda perlu mengeluarkan set hasil besar, tingkatkan nilai ini untuk menghindari pembuatan satu file berukuran terlalu besar. :param sep: Pemisah untuk file CSV. Parameter ini diabaikan untuk file non-CSV. :return: Konfigurasi untuk eksekusi SQL. """ if oss_location is None: raise ValueError("oss_location wajib diisi") if not oss_location.startswith("oss://"): raise ValueError("oss_location harus dimulai dengan oss://") if file_format != "CSV" and file_format != "PARQUET" and file_format != "ORC" and file_format != "JSON": raise ValueError("file_format harus berupa CSV, PARQUET, ORC, atau JSON") runtime_config = { # konfigurasi output sql "spark.adb.sqlOutputFormat": file_format, "spark.adb.sqlOutputPartitions": output_partitions, "spark.adb.sqlOutputLocation": oss_location, # konfigurasi csv "sep": sep } if spark_sql_runtime_config: runtime_config.update(spark_sql_runtime_config) return runtime_config def execute_sql(client: Client, dbcluster_id: str, resource_group_name: str, query: str, limit = 10000, runtime_config: dict = None, schema="default" ): """ Menjalankan pernyataan SQL di kelompok sumber daya Spark Interactive. :param client: Client Alibaba Cloud. :param dbcluster_id: ID kluster. :param resource_group_name: Kelompok sumber daya kluster. Ini harus merupakan kelompok sumber daya Spark Interactive. :param schema: Nama database default untuk eksekusi SQL. Jika Anda tidak menentukan parameter ini, nilai default akan digunakan. :param limit: Jumlah baris yang dikembalikan untuk hasil eksekusi SQL. :param query: Pernyataan SQL yang akan dijalankan. Gunakan titik koma (;) untuk memisahkan beberapa pernyataan SQL. :return: """ # Susun badan permintaan. req = ExecuteSparkWarehouseBatchSQLRequest() # ID kluster. req.dbcluster_id = dbcluster_id # Nama kelompok sumber daya. req.resource_group_name = resource_group_name # Batas waktu untuk eksekusi SQL. req.execute_time_limit_in_seconds = 3600 # Nama database tempat pernyataan SQL dijalankan. req.schema = schema # Kueri atau pernyataan SQL. req.query = query # Jumlah baris hasil yang dikembalikan. req.execute_result_limit = limit if runtime_config: # Konfigurasi untuk eksekusi SQL. req.runtime_config = json.dumps(runtime_config) # Kirim pernyataan SQL dan kembalikan ID kueri. resp: ExecuteSparkWarehouseBatchSQLResponse = client.execute_spark_warehouse_batch_sql(req) logger.info("Permintaan eksekusi kueri dikirim: {}", resp.body.data.query_id) return resp.body.data.query_id def get_query_state(client, query_id): """ Mengkueri status eksekusi pernyataan SQL. :param client: Client Alibaba Cloud. :param query_id: ID eksekusi SQL. :return: Status eksekusi dan hasil pernyataan SQL. """ req = GetSparkWarehouseBatchSQLRequest(query_id=query_id) resp: GetSparkWarehouseBatchSQLResponse = client.get_spark_warehouse_batch_sql(req) logger.info("Status kueri: {}", resp.body.data.query_state) return resp.body.data.query_state, resp def list_history_query(client, db_cluster, resource_group_name, page_num): """ Mengkueri riwayat pernyataan SQL yang dijalankan di kelompok sumber daya Spark Interactive. :param client: Client Alibaba Cloud. :param db_cluster: ID kluster. :param resource_group_name: Nama kelompok sumber daya. :param page_num: Nomor halaman untuk kueri terpaginasi. :return: Menentukan apakah masih ada halaman berikutnya. Jika kueri tersedia, Anda dapat melanjutkan ke halaman berikutnya. """ req = ListSparkWarehouseBatchSQLRequest(dbcluster_id=db_cluster, resource_group_name=resource_group_name, page_number = page_num) resp: ListSparkWarehouseBatchSQLResponse = client.list_spark_warehouse_batch_sql(req) # Jika tidak ditemukan pernyataan SQL, kembalikan True. Jika tidak, kembalikan True. Defaultnya adalah 10 entri per halaman. if resp.body.data.queries is None: return True # Cetak pernyataan SQL yang dikueri. for query in resp.body.data.queries: logger.info("ID Kueri: {}, Status: {}", query.query_id, query.query_state) logger.info("Total kueri: {}", len(resp.body.data.queries)) return len(resp.body.data.queries) < 10 def list_csv_files(oss_client, dir): for obj in oss_client.list_objects_v2(dir).object_list: if obj.key.endswith(".csv"): logger.info(f"membaca {obj.key}") # baca konten file oss csv_content = oss_client.get_object(obj.key).read().decode('utf-8') csv_reader = csv.DictReader(StringIO(csv_content)) # Cetak konten CSV for row in csv_reader: print(row) if __name__ == '__main__': logger.info("Demo ADB Spark Batch SQL") # Ganti dengan ID AccessKey Anda. _ak = "LTAI****************" # Ganti dengan rahasia AccessKey Anda. _sk = "yourAccessKeySecret" # Ganti dengan ID wilayah yang sebenarnya. _region= "cn-shanghai" # Ganti dengan ID kluster Anda. _db = "amv-uf6485635f****" # Ganti dengan nama kelompok sumber daya Anda. _rg_name = "testjob" # konfigurasi client client_config = Config( # ID AccessKey Alibaba Cloud Anda. access_key_id=_ak, # Rahasia AccessKey Alibaba Cloud Anda. access_key_secret=_sk, # Titik akhir layanan AnalyticDB for MySQL. # adb.ap-southeast-1.aliyuncs.com adalah titik akhir layanan di wilayah Tiongkok (Singapura). # adb-vpc.ap-southeast-1.aliyuncs.com digunakan dalam skenario VPC. endpoint=f"adb.{_region}.aliyuncs.com" ) # Buat client Alibaba Cloud. _client = Client(client_config) # Konfigurasi untuk eksekusi SQL. _spark_sql_runtime_config = { "spark.sql.shuffle.partitions": 1000, "spark.sql.autoBroadcastJoinThreshold": 104857600, "spark.sql.sources.partitionOverwriteMode": "dynamic", "spark.sql.sources.partitionOverwriteMode.dynamic": "dynamic" } _config = build_sql_config(oss_location="oss://testBucketName/sql_result", spark_sql_runtime_config = _spark_sql_runtime_config) # Pernyataan SQL yang akan dijalankan. _query = """ SHOW DATABASES; SELECT 100; """ _query_id = execute_sql(client = _client, dbcluster_id=_db, resource_group_name=_rg_name, query=_query, runtime_config=_config) logger.info(f"Jalankan query_id: {_query_id} untuk SQL {_query}.\n Menunggu hasil...") # Tunggu hingga eksekusi SQL selesai. current_ts = time.time() while True: query_state, resp = get_query_state(_client, _query_id) """ query_state dapat berupa salah satu dari berikut: - PENDING: Kueri sedang dalam antrian, yang dapat terjadi saat kelompok sumber daya Spark Interactive sedang memulai. - SUBMITTED: Kueri telah dikirim ke kelompok sumber daya Spark Interactive. - RUNNING: Pernyataan SQL sedang dieksekusi. - FINISHED: Eksekusi SQL berhasil diselesaikan. - FAILED: Eksekusi SQL gagal. - CANCELED: Eksekusi SQL dibatalkan. """ if query_state == "FINISHED": logger.info("kueri berhasil selesai") break elif query_state == "FAILED": # Cetak informasi kegagalan. logger.error("Info Kesalahan: {}", resp.body.data) exit(1) elif query_state == "CANCELED": # Cetak informasi pembatalan. logger.error("kueri dibatalkan") exit(1) else: time.sleep(2) if time.time() - current_ts > 600: logger.error("kueri timeout") # Jika waktu eksekusi melebihi 10 menit, batalkan eksekusi SQL. _client.cancel_spark_warehouse_batch_sql(CancelSparkWarehouseBatchSQLRequest(query_id=_query_id)) exit(1) # Satu kueri dapat berisi beberapa pernyataan. Loop berikut memproses setiap pernyataan. for stmt in resp.body.data.statements: logger.info( f"statement_id: {stmt.statement_id}, lokasi hasil: {stmt.result_uri}") # Contoh kode untuk melihat hasil. _bucket = stmt.result_uri.split("oss://")[1].split("/")[0] _dir = stmt.result_uri.replace(f"oss://{_bucket}/", "").replace("//", "/") oss_client = oss2.Bucket(oss2.Auth(client_config.access_key_id, client_config.access_key_secret), f"oss-{_region}.aliyuncs.com", _bucket) list_csv_files(oss_client, _dir) # Kueri semua pernyataan SQL yang dijalankan di kelompok sumber daya Spark Interactive. Anda dapat melakukan kueri terpaginasi. logger.info("Daftar semua riwayat kueri") page_num = 1 no_more_page = list_history_query(_client, _db, _rg_name, page_num) while no_more_page: logger.info(f"Daftar halaman {page_num}") page_num += 1 no_more_page = list_history_query(_client, _db, _rg_name, page_num)Parameter:
-
_ak: ID AccessKey akun Alibaba Cloud Anda atau pengguna RAM yang memiliki izin akses pada AnalyticDB for MySQL. Untuk informasi tentang cara mendapatkan ID AccessKey dan rahasia AccessKey, lihat Akun dan izin.
-
_sk: Rahasia AccessKey akun Alibaba Cloud Anda atau pengguna RAM yang memiliki izin akses pada AnalyticDB for MySQL. Untuk informasi tentang cara mendapatkan ID AccessKey dan rahasia AccessKey, lihat Akun dan izin.
-
_region: ID wilayah tempat kluster AnalyticDB for MySQL Anda berada.
-
_db: ID kluster AnalyticDB for MySQL.
-
_rg_name: Nama kelompok sumber daya Spark Interactive.
-
oss_location (opsional): Jalur OSS tempat file hasil kueri disimpan.
Jika Anda tidak menentukan parameter ini, Anda hanya dapat melihat lima baris pertama hasil kueri di Log pada halaman .
-
Aplikasi
Hive JDBC
-
Di file pom.xml, konfigurasikan dependensi Maven.
<dependency> <groupId>org.apache.hive</groupId> <artifactId>hive-jdbc</artifactId> <version>2.3.9</version> </dependency> -
Buat koneksi dan jalankan pernyataan Spark SQL.
public class java { public static void main(String[] args) throws Exception { Class.forName("org.apache.hive.jdbc.HiveDriver"); String url = "<JDBC-connection-string>"; Connection con = DriverManager.getConnection(url, "<username>", "<password>"); Statement stmt = con.createStatement(); ResultSet tables = stmt.executeQuery("show tables"); List<String> tbls = new ArrayList<>(); while (tables.next()) { System.out.println(tables.getString("tableName")); tbls.add(tables.getString("tableName")); } } }Parameter:
-
String koneksi JDBC: String koneksi JDBC kelompok sumber daya Spark Interactive yang Anda peroleh di bagian Persiapan. Ganti
defaultdengan nama database yang ingin Anda hubungkan. -
Username: AnalyticDB for MySQL untuk AnalyticDB for MySQL.
-
Password: Kata sandi AnalyticDB for MySQL untuk AnalyticDB for MySQL.
-
PyHive
-
Instal client PyHive.
pip install pyhive -
Buat koneksi dan jalankan pernyataan Spark SQL.
from pyhive import hive from TCLIService.ttypes import TOperationState cursor = hive.connect( host='<endpoint>', port=<port>, username='<resource_group_name>/<username>', password='<password>', auth='CUSTOM' ).cursor() cursor.execute('show tables') status = cursor.poll().operationState while status in (TOperationState.INITIALIZED_STATE, TOperationState.RUNNING_STATE): logs = cursor.fetch_logs() for message in logs: print(message) # Jika diperlukan, kueri asinkron dapat dibatalkan kapan saja dengan: # cursor.cancel() status = cursor.poll().operationState print(cursor.fetchall())Parameter:
-
Endpoint: Titik akhir kelompok sumber daya Spark Interactive yang Anda peroleh di bagian Persiapan.
-
Port: Port kelompok sumber daya Spark Interactive, yaitu 10000.
-
Nama kelompok sumber daya: Nama kelompok sumber daya Spark Interactive.
-
Username: AnalyticDB for MySQL untuk AnalyticDB for MySQL.
-
Password: Kata sandi AnalyticDB for MySQL untuk AnalyticDB for MySQL.
-
Client
Selain client Beeline, DBeaver, DBVisualizer, dan DataGrip yang dijelaskan dalam topik ini, Anda juga dapat melakukan analisis interaktif di tool penjadwalan alur kerja seperti Airflow, Azkaban, dan DolphinScheduler.
Beeline
-
Hubungkan ke kelompok sumber daya Spark Interactive.
Gunakan format perintah berikut:
!connect <JDBC-connection-string> <username> <password>-
String koneksi JDBC: String koneksi JDBC kelompok sumber daya Spark Interactive yang Anda peroleh di bagian Persiapan. Ganti
defaultdengan nama database yang ingin Anda hubungkan. -
Username: AnalyticDB for MySQL untuk AnalyticDB for MySQL.
-
Password: Kata sandi AnalyticDB for MySQL untuk AnalyticDB for MySQL.
Contoh:
!connect jdbc:hive2://amv-bp1c3em7b2e****-spark.ads.aliyuncs.com:10000/adb_test spark_resourcegroup/AdbSpark14**** Spark23****Koneksi yang berhasil akan menghasilkan output berikut:
Connected to: Spark SQL (version 3.2.0) Driver: Hive JDBC (version 2.3.9) Transaction isolation: TRANSACTION_REPEATABLE_READ -
-
Jalankan pernyataan Spark SQL.
SHOW TABLES;
DBeaver
-
Buka client DBeaver dan pilih .
-
Pada halaman Connect to a database, pilih Apache Spark lalu klik Next.
-
Konfigurasikan Hadoop/Apache Spark connection settings sebagai berikut:
Parameter
Deskripsi
Metode koneksi
Pilih URL.
URL JDBC
Masukkan string koneksi JDBC yang Anda peroleh di bagian Persiapan.
PentingGanti
defaultdalam string koneksi dengan nama database Anda.Username
AnalyticDB for MySQL untuk AnalyticDB for MySQL.
Password
Kata sandi AnalyticDB for MySQL untuk AnalyticDB for MySQL.
-
Setelah mengonfigurasi parameter, klik Test Connection.
PentingPertama kali menguji koneksi, DBeaver akan meminta Anda mengunduh driver yang diperlukan. Klik Download untuk mengunduhnya.
-
Setelah pengujian koneksi berhasil, klik Finish.
-
Pada tab Database Navigator, perluas sumber data lalu klik database.
-
Di editor kode di sebelah kanan, masukkan pernyataan SQL lalu klik ikon
untuk menjalankannya.SHOW TABLES;+-----------+-----------+-------------+ | namespace | tableName | isTemporary | +-----------+-----------+-------------+ | db | test | [] | +-----------+-----------+-------------+
DBVisualizer
-
Buka client DBVisualizer dan pilih .
-
Pada halaman Driver Manager, pilih Hive lalu klik ikon
. -
Pada tab Driver Settings, konfigurasikan parameter berikut:
Parameter
Deskripsi
Name
Nama kustom untuk sumber data Hive.
Format URL
Masukkan string koneksi JDBC yang Anda peroleh di bagian Persiapan.
PentingGanti
defaultdalam string koneksi dengan nama database Anda.Kelas driver
Pilih org.apache.hive.jdbc.HiveDriver.
PentingSetelah mengonfigurasi parameter, klik Start Download untuk mengunduh driver.
-
Setelah driver diunduh, pilih .
-
Pada kotak dialog Create Database Connection from Database URL, konfigurasikan parameter yang dijelaskan dalam tabel berikut.
Parameter
Deskripsi
URL Basis Data
Masukkan string koneksi JDBC yang Anda peroleh di bagian Persiapan.
PentingGanti
defaultdalam string koneksi dengan nama database Anda.Kelas driver
Pilih sumber data Hive yang Anda buat di Langkah 3.
-
Pada halaman Connection, konfigurasikan parameter koneksi berikut lalu klik Connect.
Parameter
Deskripsi
Name
Secara default, parameter ini diatur ke nama sumber data Hive yang Anda buat di Langkah 3. Anda dapat menyesuaikan nama tersebut.
Notes
Masukkan catatan.
Jenis driver
Pilih Hive.
URL Database
Masukkan string koneksi JDBC yang Anda peroleh di bagian Persiapan.
PentingGanti
defaultdalam string koneksi dengan nama database Anda.Database userid
AnalyticDB for MySQL untuk AnalyticDB for MySQL.
Database password
Kata sandi AnalyticDB for MySQL untuk AnalyticDB for MySQL.
CatatanBiarkan parameter lainnya pada nilai default.
-
Setelah koneksi terbentuk, pada tab Database, perluas sumber data lalu klik database.
-
Di editor kode di sebelah kanan, masukkan pernyataan SQL lalu klik ikon
untuk menjalankannya.SHOW TABLES;+-----------+-----------+-------------+ | namespace | tableName | isTemporary | +-----------+-----------+-------------+ | db | test | false | +-----------+-----------+-------------+
DataGrip
-
Buka client DataGrip, pilih , lalu buat proyek.
-
Tambahkan sumber data.
-
Klik ikon
lalu pilih . -
Pada kotak dialog Data Sources and Drivers yang muncul, konfigurasikan parameter berikut lalu klik OK.
Atur Driver ke Apache Spark dan Authentication ke User & Password.
Parameter
Deskripsi
Name
Nama sumber data, yang dapat Anda sesuaikan. Topik ini menggunakan
adbtestsebagai contoh.Host
Masukkan string koneksi JDBC yang Anda peroleh di bagian Persiapan.
PentingGanti
defaultdalam string koneksi dengan nama database Anda.Port
Port kelompok sumber daya Spark Interactive, yaitu 10000.
User
AnalyticDB for MySQL untuk AnalyticDB for MySQL.
Password
Kata sandi AnalyticDB for MySQL untuk AnalyticDB for MySQL.
Schema
Nama database di kluster AnalyticDB for MySQL.
-
-
Jalankan pernyataan Spark SQL.
-
Di daftar sumber data, klik kanan sumber data yang Anda buat di Langkah 2 lalu pilih .
-
Di panel Console yang muncul, jalankan pernyataan Spark SQL.
SHOW TABLES;
-
Alat BI
Anda dapat melakukan analisis interaktif di tool BI seperti Redash, Power BI, dan Metabase.