يستفيد هذا الموصل من تحسينات خاصة بـ ClickHouse، مثل التقسيم المتقدم ودفع شروط التصفية، بهدف تحسين أداء الاستعلامات ومعالجة البيانات. ويعتمد هذا الموصل على موصل JDBC الرسمي لـ ClickHouse، كما يدير الكتالوج الخاص به.
قبل Spark 3.0، لم يكن Spark يوفّر مفهوم كتالوج مضمّن، لذلك كان المستخدمون يعتمدون عادةً على أنظمة كتالوج خارجية مثل Hive Metastore أو AWS Glue. ومع هذه الحلول الخارجية، كان على المستخدمين تسجيل جداول مصدر البيانات الخاصة بهم يدويًا قبل أن يتمكنوا من الوصول إليها في Spark. لكن منذ أن قدّم Spark 3.0 مفهوم الكتالوج، أصبح بإمكان Spark الآن اكتشاف الجداول تلقائيًا من خلال تسجيل إضافات الكتالوج.
الكتالوج الافتراضي في Spark هو spark_catalog، وتُعرَّف الجداول بالشكل {catalog name}.{database}.{table}. ومع
ميزة الكتالوج الجديدة، أصبح من الممكن الآن إضافة عدة كتالوجات والعمل بها ضمن تطبيق Spark واحد.
المفاضلة بين Catalog API وTableProvider API
يدعم موصل ClickHouse Spark نمطين للوصول: Catalog API وTableProvider API (الوصول المستند إلى التنسيق). ويساعدك فهم الفروق بينهما على اختيار النهج الأنسب لحالة الاستخدام الخاصة بك.
Catalog API مقابل TableProvider API
| الميزة | Catalog API | TableProvider API |
|---|---|---|
| الإعداد | مركزي عبر إعدادات Spark | لكل عملية عبر الخيارات |
| اكتشاف الجداول | تلقائي عبر كتالوج | تحديد الجدول يدويًا |
| عمليات DDL | دعم كامل (CREATE, DROP, ALTER) | محدود (إنشاء الجدول تلقائيًا فقط) |
| تكامل Spark SQL | أصلي (clickhouse.database.table) |
يتطلب تحديد التنسيق |
| حالة الاستخدام | اتصالات طويلة الأمد ومستقرة مع إعداد مركزي | وصول مخصص أو ديناميكي أو مؤقت |
المتطلبات
- Java 8 أو 17 (يتطلب Spark 4.0 Java 17 أو أحدث)
- Scala 2.12 أو 2.13 (يدعم Spark 4.0 Scala 2.13 فقط)
- Apache Spark 3.3 أو 3.4 أو 3.5 أو 4.0
مصفوفة التوافق
| الإصدار | إصدارات Spark المتوافقة | إصدار ClickHouse JDBC |
|---|---|---|
| main | Spark 3.3, 3.4, 3.5, 4.0 | 0.9.4 |
| 0.10.0 | Spark 3.3, 3.4, 3.5, 4.0 | 0.9.5 |
| 0.9.0 | Spark 3.3, 3.4, 3.5, 4.0 | 0.9.4 |
| 0.8.1 | Spark 3.3, 3.4, 3.5 | 0.6.3 |
| 0.7.3 | Spark 3.3, 3.4 | 0.4.6 |
| 0.6.0 | Spark 3.3 | 0.3.2-patch11 |
| 0.5.0 | Spark 3.2, 3.3 | 0.3.2-patch11 |
| 0.4.0 | Spark 3.2, 3.3 | لا يعتمد عليه |
| 0.3.0 | Spark 3.2, 3.3 | لا يعتمد عليه |
| 0.2.1 | Spark 3.2 | لا يعتمد عليه |
| 0.1.2 | Spark 3.2 | لا يعتمد عليه |
التثبيت والإعداد
لدمج ClickHouse مع Spark، تتوفر عدة خيارات للتثبيت تناسب إعدادات المشاريع المختلفة.
يمكنك إضافة موصل ClickHouse Spark كتبعية مباشرةً في ملف البناء الخاص بمشروعك (مثل pom.xml
في Maven أو build.sbt في SBT).
وبدلاً من ذلك، يمكنك وضع ملفات JAR المطلوبة في المجلد $SPARK_HOME/jars/، أو تمريرها مباشرةً إلى Spark
باستخدام الخيار --jars في الأمر spark-submit.
ويضمن كلا النهجين توفّر موصل ClickHouse في بيئة Spark لديك.
الاستيراد كتبعية
<dependency>
<groupId>com.clickhouse.spark</groupId>
<artifactId>clickhouse-spark-runtime-{{ spark_binary_version }}_{{ scala_binary_version }}</artifactId>
<version>{{ stable_version }}</version>
</dependency>
<dependency>
<groupId>com.clickhouse</groupId>
<artifactId>clickhouse-jdbc</artifactId>
<classifier>all</classifier>
<version>{{ clickhouse_jdbc_version }}</version>
<exclusions>
<exclusion>
<groupId>*</groupId>
<artifactId>*</artifactId>
</exclusion>
</exclusions>
</dependency>لاستخدام إصدار SNAPSHOT، اتبع تعليمات Sonatype لاستخدام إصدارات SNAPSHOT مع Maven.
dependencies {
implementation("com.clickhouse.spark:clickhouse-spark-runtime-{{ spark_binary_version }}_{{ scala_binary_version }}:{{ stable_version }}")
implementation("com.clickhouse:clickhouse-jdbc:{{ clickhouse_jdbc_version }}:all") { transitive = false }
}لاستخدام إصدار SNAPSHOT، اتبع تعليمات Sonatype لاستخدام إصدارات SNAPSHOT مع Gradle.
libraryDependencies += "com.clickhouse" % "clickhouse-jdbc" % {{ clickhouse_jdbc_version }} classifier "all"
libraryDependencies += "com.clickhouse.spark" %% clickhouse-spark-runtime-{{ spark_binary_version }}_{{ scala_binary_version }} % {{ stable_version }}عند العمل مع خيارات shell في Spark (Spark SQL CLI وSpark Shell CLI وأمر Spark Submit)، يمكن تسجيل التبعيات عبر تمرير ملفات JARs المطلوبة:
$SPARK_HOME/bin/spark-sql \
--jars /path/clickhouse-spark-runtime-{{ spark_binary_version }}_{{ scala_binary_version }}:{{ stable_version }}.jar,/path/clickhouse-jdbc-{{ clickhouse_jdbc_version }}-all.jarإذا كنت تريد تجنب نسخ ملفات JARs إلى عقدة عميل Spark، يمكنك استخدام ما يلي بدلًا من ذلك:
--repositories https://{maven-central-mirror or private-nexus-repo} \
--packages com.clickhouse.spark:clickhouse-spark-runtime-{{ spark_binary_version }}_{{ scala_binary_version }}:{{ stable_version }},com.clickhouse:clickhouse-jdbc:{{ clickhouse_jdbc_version }}ملاحظة: بالنسبة إلى حالات الاستخدام التي تقتصر على SQL فقط، يُوصى باستخدام Apache Kyuubi في بيئات production.
نزّل المكتبة
نمط تسمية ملف JAR التنفيذي هو:
clickhouse-spark-runtime-${spark_binary_version}_${scala_binary_version}-${version}.jarيمكنك العثور على جميع ملفات JAR المتاحة للإصدارات المنشورة في Maven Central Repository. ملفات JAR اليومية من نوع SNAPSHOT متاحة عبر مستودع Sonatype للقطات snapshots المُعدّ أعلاه.
سجّل الكتالوج (مطلوب)
للوصول إلى جداول ClickHouse الخاصة بك، يجب تهيئة كتالوج Spark جديد باستخدام إعدادات config التالية:
| الخاصية | القيمة | القيمة الافتراضية | مطلوب |
|---|---|---|---|
spark.sql.catalog.<catalog_name> |
com.clickhouse.spark.ClickHouseCatalog |
N/A |
نعم |
spark.sql.catalog.<catalog_name>.host |
<clickhouse_host> |
localhost |
لا |
spark.sql.catalog.<catalog_name>.protocol |
http |
http |
لا |
spark.sql.catalog.<catalog_name>.http_port |
<clickhouse_port> |
8123 |
لا |
spark.sql.catalog.<catalog_name>.user |
<clickhouse_username> |
default |
لا |
spark.sql.catalog.<catalog_name>.password |
<clickhouse_password> |
(سلسلة فارغة) | لا |
spark.sql.catalog.<catalog_name>.database |
<database> |
default |
لا |
spark.<catalog_name>.write.format |
json |
arrow |
لا |
يمكن ضبط هذه الإعدادات بإحدى الطرق التالية:
- تعديل/إنشاء
spark-defaults.conf. - تمرير الإعدادات إلى الأمر
spark-submit(أو إلى أوامر CLI spark-shell/spark-sql). - إضافة الإعدادات عند تهيئة السياق.
استخدام واجهة برمجة تطبيقات TableProvider (الوصول المستند إلى التنسيق)
إلى جانب النهج المستند إلى كتالوج، يدعم موصل ClickHouse Spark نمط وصول مستندًا إلى التنسيق عبر واجهة برمجة تطبيقات TableProvider.
مثال على القراءة بأسلوب مستند إلى التنسيق
from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()
# القراءة من ClickHouse باستخدام واجهة برمجة تطبيقات format
df = spark.read \
.format("clickhouse") \
.option("host", "your-clickhouse-host") \
.option("protocol", "https") \
.option("http_port", "8443") \
.option("database", "default") \
.option("table", "your_table") \
.option("user", "default") \
.option("password", "your_password") \
.option("ssl", "true") \
.load()
df.show()val df = spark.read
.format("clickhouse")
.option("host", "your-clickhouse-host")
.option("protocol", "https")
.option("http_port", "8443")
.option("database", "default")
.option("table", "your_table")
.option("user", "default")
.option("password", "your_password")
.option("ssl", "true")
.load()
df.show()Dataset<Row> df = spark.read()
.format("clickhouse")
.option("host", "your-clickhouse-host")
.option("protocol", "https")
.option("http_port", "8443")
.option("database", "default")
.option("table", "your_table")
.option("user", "default")
.option("password", "your_password")
.option("ssl", "true")
.load();
df.show();مثال على الكتابة المستندة إلى التنسيق
# الكتابة إلى ClickHouse باستخدام واجهة برمجة تطبيقات format
df.write \
.format("clickhouse") \
.option("host", "your-clickhouse-host") \
.option("protocol", "https") \
.option("http_port", "8443") \
.option("database", "default") \
.option("table", "your_table") \
.option("user", "default") \
.option("password", "your_password") \
.option("ssl", "true") \
.mode("append") \
.save()df.write
.format("clickhouse")
.option("host", "your-clickhouse-host")
.option("protocol", "https")
.option("http_port", "8443")
.option("database", "default")
.option("table", "your_table")
.option("user", "default")
.option("password", "your_password")
.option("ssl", "true")
.mode("append")
.save()df.write()
.format("clickhouse")
.option("host", "your-clickhouse-host")
.option("protocol", "https")
.option("http_port", "8443")
.option("database", "default")
.option("table", "your_table")
.option("user", "default")
.option("password", "your_password")
.option("ssl", "true")
.mode("append")
.save();ميزات TableProvider API
توفّر TableProvider API عدة ميزات قوية:
الإنشاء التلقائي للجدول
عند الكتابة إلى جدول غير موجود، ينشئ الموصل الجدول تلقائيًا ببنية مناسبة. ويوفر الموصل إعدادات افتراضية ذكية:
- المحرّك: تكون القيمة الافتراضية
MergeTree()إذا لم يتم تحديده. يمكنك تحديد محرّك مختلف باستخدام الخيارengine(مثلReplacingMergeTree(),SummingMergeTree(), وغيرها) - ORDER BY: مطلوب - يجب تحديد الخيار
order_byصراحةً عند إنشاء جدول جديد. ويتحقق الموصل من أن جميع الأعمدة المحددة موجودة في البنية. - دعم المفاتيح القابلة لـ NULL: يضيف تلقائيًا
settings.allow_nullable_key=1إذا كانت عبارة ORDER BY تحتوي على أعمدة قابلة لـ NULL
# سيتم إنشاء الجدول تلقائيًا مع تحديد ORDER BY صراحةً (مطلوب)
df.write \
.format("clickhouse") \
.option("host", "your-host") \
.option("database", "default") \
.option("table", "new_table") \
.option("order_by", "id") \
.mode("append") \
.save()
# تحديد خيارات إنشاء الجدول باستخدام محرّك مخصص
df.write \
.format("clickhouse") \
.option("host", "your-host") \
.option("database", "default") \
.option("table", "new_table") \
.option("order_by", "id, timestamp") \
.option("engine", "ReplacingMergeTree()") \
.option("settings.allow_nullable_key", "1") \
.mode("append") \
.save()// سيتم إنشاء الجدول تلقائيًا مع تحديد ORDER BY صراحةً (مطلوب)
df.write
.format("clickhouse")
.option("host", "your-host")
.option("database", "default")
.option("table", "new_table")
.option("order_by", "id")
.mode("append")
.save()
// مع تحديد خيارات إنشاء الجدول صراحةً واستخدام محرّك مخصص
df.write
.format("clickhouse")
.option("host", "your-host")
.option("database", "default")
.option("table", "new_table")
.option("order_by", "id")
.mode("append")
.save()
// مع تحديد خيارات إنشاء الجدول صراحةً واستخدام محرّك مخصص
df.write
.format("clickhouse")
.option("host", "your-host")
.option("database", "default")
.option("table", "new_table")
.option("order_by", "id, timestamp")
.option("engine", "ReplacingMergeTree()")
.option("settings.allow_nullable_key", "1")
.mode("append")
.save()// سيتم إنشاء الجدول تلقائيًا مع تحديد ORDER BY صراحةً (مطلوب)
df.write()
.format("clickhouse")
.option("host", "your-host")
.option("database", "default")
.option("table", "new_table")
.option("order_by", "id")
.mode("append")
.save();
// مع تحديد خيارات إنشاء الجدول صراحةً واستخدام محرّك مخصص
df.write()
.format("clickhouse")
.option("host", "your-host")
.option("database", "default")
.option("table", "new_table")
.option("order_by", "id, timestamp")
.option("engine", "ReplacingMergeTree()")
.option("settings.allow_nullable_key", "1")
.mode("append")
.save();خيارات اتصال TableProvider
عند استخدام واجهة برمجة التطبيقات المعتمدة على التنسيق، تتوفر خيارات الاتصال التالية:
خيارات الاتصال
| الخيار | الوصف | القيمة الافتراضية | مطلوب |
|---|---|---|---|
host |
اسم مضيف خادم ClickHouse | localhost |
نعم |
protocol |
بروتوكول الاتصال (http أو https) |
http |
لا |
http_port |
منفذ HTTP/HTTPS | 8123 |
لا |
database |
اسم قاعدة البيانات | default |
نعم |
table |
اسم الجدول | N/A | نعم |
user |
اسم المستخدم لأغراض authentication | default |
لا |
password |
كلمة المرور لأغراض authentication | (سلسلة فارغة) | لا |
ssl |
تمكين اتصال SSL | false |
لا |
ssl_mode |
وضع SSL (NONE، STRICT، إلخ) |
STRICT |
لا |
timezone |
المنطقة الزمنية لعمليات التاريخ/الوقت | server |
لا |
خيارات إنشاء الجدول
تُستخدم هذه الخيارات عندما لا يكون الجدول موجودًا ويحتاج إلى إنشائه:
| الخيار | الوصف | القيمة الافتراضية | مطلوب |
|---|---|---|---|
order_by |
الأعمدة المستخدمة في عبارة ORDER BY. تُفصل بفواصل عند وجود عدة أعمدة | N/A | نعم |
engine |
محرك جدول ClickHouse (مثل MergeTree(), ReplacingMergeTree(), SummingMergeTree()، إلخ) |
MergeTree() |
لا |
settings.allow_nullable_key |
تمكين المفاتيح التي تقبل NULL في ORDER BY (لـ ClickHouse Cloud) | يُكتشف تلقائيًا** | لا |
settings.<key> |
أي إعداد لجدول ClickHouse | N/A | لا |
cluster |
اسم العنقود لجداول Distributed | N/A | لا |
clickhouse.column.<name>.variant_types |
قائمة مفصولة بفواصل بأنواع ClickHouse لأعمدة Variant (مثل String, Int64, Bool, JSON). أسماء الأنواع حساسة لحالة الأحرف. المسافات بعد الفواصل اختيارية. |
N/A | لا |
- الخيار
order_byمطلوب عند إنشاء جدول جديد. يجب أن تكون جميع الأعمدة المحددة موجودة في المخطط. ** يُضبط تلقائيًا على1إذا كان ORDER BY يحتوي على أعمدة تقبل NULL ولم يُحدَّد هذا الإعداد صراحةً.
أوضاع الكتابة
يدعم Spark connector (كلًّا من TableProvider API وCatalog API) أوضاع الكتابة التالية في Spark:
append: إضافة البيانات إلى جدول موجودoverwrite: استبدال جميع البيانات في الجدول (مع تفريغ الجدول)
# وضع الاستبدال (يُفرِّغ الجدول أولًا)
df.write \
.format("clickhouse") \
.option("host", "your-host") \
.option("database", "default") \
.option("table", "my_table") \
.mode("overwrite") \
.save()// وضع الاستبدال (يُفرِّغ الجدول أولًا)
df.write
.format("clickhouse")
.option("host", "your-host")
.option("database", "default")
.option("table", "my_table")
.mode("overwrite")
.save()// وضع الاستبدال (يُفرِّغ الجدول أولًا)
df.write()
.format("clickhouse")
.option("host", "your-host")
.option("database", "default")
.option("table", "my_table")
.mode("overwrite")
.save();تهيئة خيارات ClickHouse
يدعم كلٌّ من Catalog API وTableProvider API تهيئة الخيارات الخاصة بـ ClickHouse (وليس خيارات الموصل). وتُمرَّر هذه الخيارات إلى ClickHouse عند إنشاء الجداول أو تنفيذ الاستعلامات.
تتيح لك خيارات ClickHouse ضبط إعدادات خاصة به مثل allow_nullable_key وindex_granularity، بالإضافة إلى إعدادات أخرى على مستوى الجدول أو على مستوى الاستعلام. وهي تختلف عن خيارات الموصل (مثل host وdatabase وtable) التي تتحكم في كيفية اتصال الموصل بـ ClickHouse.
استخدام واجهة TableProvider API
مع واجهة TableProvider API، استخدم تنسيق الخيار settings.<key>:
df.write \
.format("clickhouse") \
.option("host", "your-host") \
.option("database", "default") \
.option("table", "my_table") \
.option("order_by", "id") \
.option("settings.allow_nullable_key", "1") \
.option("settings.index_granularity", "8192") \
.mode("append") \
.save()df.write
.format("clickhouse")
.option("host", "your-host")
.option("database", "default")
.option("table", "my_table")
.option("order_by", "id")
.option("settings.allow_nullable_key", "1")
.option("settings.index_granularity", "8192")
.mode("append")
.save()df.write()
.format("clickhouse")
.option("host", "your-host")
.option("database", "default")
.option("table", "my_table")
.option("order_by", "id")
.option("settings.allow_nullable_key", "1")
.option("settings.index_granularity", "8192")
.mode("append")
.save();استخدام Catalog API
باستخدام Catalog API، استخدم الصيغة spark.sql.catalog.<catalog_name>.option.<key> في إعدادات Spark:
spark.sql.catalog.clickhouse.option.allow_nullable_key 1
spark.sql.catalog.clickhouse.option.index_granularity 8192أو اضبطها عند إنشاء الجداول باستخدام Spark SQL:
CREATE TABLE clickhouse.default.my_table (
id INT,
name STRING
) USING ClickHouse
TBLPROPERTIES (
engine = 'MergeTree()',
order_by = 'id',
'settings.allow_nullable_key' = '1',
'settings.index_granularity' = '8192'
)إعدادات ClickHouse Cloud
عند الاتصال بـ ClickHouse Cloud، تأكد من تمكين SSL وضبط وضع SSL المناسب. على سبيل المثال:
spark.sql.catalog.clickhouse.option.ssl true
spark.sql.catalog.clickhouse.option.ssl_mode NONEقراءة البيانات
public static void main(String[] args) {
// إنشاء جلسة Spark
SparkSession spark = SparkSession.builder()
.appName("example")
.master("local[*]")
.config("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
.config("spark.sql.catalog.clickhouse.host", "127.0.0.1")
.config("spark.sql.catalog.clickhouse.protocol", "http")
.config("spark.sql.catalog.clickhouse.http_port", "8123")
.config("spark.sql.catalog.clickhouse.user", "default")
.config("spark.sql.catalog.clickhouse.password", "123456")
.config("spark.sql.catalog.clickhouse.database", "default")
.config("spark.clickhouse.write.format", "json")
.getOrCreate();
Dataset<Row> df = spark.sql("select * from clickhouse.default.example_table");
df.show();
spark.stop();
}object NativeSparkRead extends App {
val spark = SparkSession.builder
.appName("example")
.master("local[*]")
.config("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
.config("spark.sql.catalog.clickhouse.host", "127.0.0.1")
.config("spark.sql.catalog.clickhouse.protocol", "http")
.config("spark.sql.catalog.clickhouse.http_port", "8123")
.config("spark.sql.catalog.clickhouse.user", "default")
.config("spark.sql.catalog.clickhouse.password", "123456")
.config("spark.sql.catalog.clickhouse.database", "default")
.config("spark.clickhouse.write.format", "json")
.getOrCreate
val df = spark.sql("select * from clickhouse.default.example_table")
df.show()
spark.stop()
}from pyspark.sql import SparkSession
packages = [
"com.clickhouse.spark:clickhouse-spark-runtime-3.4_2.12:0.8.0",
"com.clickhouse:clickhouse-client:0.7.0",
"com.clickhouse:clickhouse-http-client:0.7.0",
"org.apache.httpcomponents.client5:httpclient5:5.2.1"
]
spark = (SparkSession.builder
.config("spark.jars.packages", ",".join(packages))
.getOrCreate())
spark.conf.set("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
spark.conf.set("spark.sql.catalog.clickhouse.host", "127.0.0.1")
spark.conf.set("spark.sql.catalog.clickhouse.protocol", "http")
spark.conf.set("spark.sql.catalog.clickhouse.http_port", "8123")
spark.conf.set("spark.sql.catalog.clickhouse.user", "default")
spark.conf.set("spark.sql.catalog.clickhouse.password", "123456")
spark.conf.set("spark.sql.catalog.clickhouse.database", "default")
spark.conf.set("spark.clickhouse.write.format", "json")
df = spark.sql("select * from clickhouse.default.example_table")
df.show() CREATE TEMPORARY VIEW jdbcTable
USING org.apache.spark.sql.jdbc
OPTIONS (
url "jdbc:ch://localhost:8123/default",
dbtable "schema.tablename",
user "username",
password "password",
driver "com.clickhouse.jdbc.ClickHouseDriver"
);
SELECT * FROM jdbcTable;كتابة البيانات
public static void main(String[] args) throws AnalysisException {
// إنشاء جلسة Spark
SparkSession spark = SparkSession.builder()
.appName("example")
.master("local[*]")
.config("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
.config("spark.sql.catalog.clickhouse.host", "127.0.0.1")
.config("spark.sql.catalog.clickhouse.protocol", "http")
.config("spark.sql.catalog.clickhouse.http_port", "8123")
.config("spark.sql.catalog.clickhouse.user", "default")
.config("spark.sql.catalog.clickhouse.password", "123456")
.config("spark.sql.catalog.clickhouse.database", "default")
.config("spark.clickhouse.write.format", "json")
.getOrCreate();
// تحديد schema الخاصة بـ DataFrame
StructType schema = new StructType(new StructField[]{
DataTypes.createStructField("id", DataTypes.IntegerType, false),
DataTypes.createStructField("name", DataTypes.StringType, false),
});
List<Row> data = Arrays.asList(
RowFactory.create(1, "Alice"),
RowFactory.create(2, "Bob")
);
// إنشاء DataFrame
Dataset<Row> df = spark.createDataFrame(data, schema);
df.writeTo("clickhouse.default.example_table").append();
spark.stop();
}object NativeSparkWrite extends App {
// إنشاء جلسة Spark
val spark: SparkSession = SparkSession.builder
.appName("example")
.master("local[*]")
.config("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
.config("spark.sql.catalog.clickhouse.host", "127.0.0.1")
.config("spark.sql.catalog.clickhouse.protocol", "http")
.config("spark.sql.catalog.clickhouse.http_port", "8123")
.config("spark.sql.catalog.clickhouse.user", "default")
.config("spark.sql.catalog.clickhouse.password", "123456")
.config("spark.sql.catalog.clickhouse.database", "default")
.config("spark.clickhouse.write.format", "json")
.getOrCreate
// تحديد schema الخاصة بـ DataFrame
val rows = Seq(Row(1, "John"), Row(2, "Doe"))
val schema = List(
StructField("id", DataTypes.IntegerType, nullable = false),
StructField("name", StringType, nullable = true)
)
// إنشاء df
val df: DataFrame = spark.createDataFrame(
spark.sparkContext.parallelize(rows),
StructType(schema)
)
df.writeTo("clickhouse.default.example_table").append()
spark.stop()
}from pyspark.sql import SparkSession
from pyspark.sql import Row
# يمكنك استخدام أي مجموعة أخرى من الحزم ما دامت تستوفي Compatibility Matrix المذكورة أعلاه.
packages = [
"com.clickhouse.spark:clickhouse-spark-runtime-3.4_2.12:0.8.0",
"com.clickhouse:clickhouse-client:0.7.0",
"com.clickhouse:clickhouse-http-client:0.7.0",
"org.apache.httpcomponents.client5:httpclient5:5.2.1"
]
spark = (SparkSession.builder
.config("spark.jars.packages", ",".join(packages))
.getOrCreate())
spark.conf.set("spark.sql.catalog.clickhouse", "com.clickhouse.spark.ClickHouseCatalog")
spark.conf.set("spark.sql.catalog.clickhouse.host", "127.0.0.1")
spark.conf.set("spark.sql.catalog.clickhouse.protocol", "http")
spark.conf.set("spark.sql.catalog.clickhouse.http_port", "8123")
spark.conf.set("spark.sql.catalog.clickhouse.user", "default")
spark.conf.set("spark.sql.catalog.clickhouse.password", "123456")
spark.conf.set("spark.sql.catalog.clickhouse.database", "default")
spark.conf.set("spark.clickhouse.write.format", "json")
# إنشاء DataFrame
data = [Row(id=11, name="John"), Row(id=12, name="Doe")]
df = spark.createDataFrame(data)
# كتابة DataFrame إلى ClickHouse
df.writeTo("clickhouse.default.example_table").append() -- resultTable هو DataFrame الوسيط في Spark الذي نريد إدراجه في clickhouse.default.example_table
INSERT INTO TABLE clickhouse.default.example_table
SELECT * FROM resultTable;
عمليات DDL
يمكنك تنفيذ عمليات DDL على مثيل ClickHouse لديك باستخدام Spark SQL، مع حفظ جميع التغييرات فورًا في ClickHouse. يتيح لك Spark SQL كتابة الاستعلامات تمامًا كما تفعل في ClickHouse، لذا يمكنك تنفيذ أوامر مثل CREATE TABLE وTRUNCATE وغيرها مباشرةً، من دون أي تعديل، على سبيل المثال:
USE clickhouse;
CREATE TABLE test_db.tbl_sql (
create_time TIMESTAMP NOT NULL,
m INT NOT NULL COMMENT 'part key',
id BIGINT NOT NULL COMMENT 'sort key',
value STRING
) USING ClickHouse
PARTITIONED BY (m)
TBLPROPERTIES (
engine = 'MergeTree()',
order_by = 'id',
settings.index_granularity = 8192
);توضح الأمثلة أعلاه استعلامات Spark SQL، ويمكنك تشغيلها داخل تطبيقك باستخدام أي واجهة برمجة تطبيقات، مثل Java أو Scala أو PySpark أو shell.
العمل مع VariantType
يدعم الموصل VariantType في Spark للعمل مع البيانات شبه المهيكلة. ويرتبط VariantType بالنوعين JSON وVariant في ClickHouse، ما يتيح لك تخزين البيانات ذات المخطط المرن والاستعلام عنها بكفاءة.
تعيين أنواع ClickHouse
| نوع ClickHouse | نوع Spark | الوصف |
|---|---|---|
JSON |
VariantType |
يخزّن كائنات JSON فقط (يجب أن يبدأ بـ {) |
Variant(T1, T2, ...) |
VariantType |
يخزّن أنواعًا متعددة، بما في ذلك الأنواع الأولية والمصفوفات وJSON |
قراءة بيانات VariantType
عند القراءة من ClickHouse، تُعيَّن الأعمدة JSON وVariant تلقائيًا إلى VariantType في Spark:
// قراءة عمود JSON كـ VariantType
val df = spark.sql("SELECT id, data FROM clickhouse.default.json_table")
// الوصول إلى بيانات Variant
df.show()
// تحويل Variant إلى سلسلة JSON للفحص
import org.apache.spark.sql.functions._
df.select(
col("id"),
to_json(col("data")).as("data_json")
).show()# قراءة عمود JSON كـ VariantType
df = spark.sql("SELECT id, data FROM clickhouse.default.json_table")
# الوصول إلى بيانات Variant
df.show()
# تحويل Variant إلى سلسلة JSON للفحص
from pyspark.sql.functions import to_json
df.select(
"id",
to_json("data").alias("data_json")
).show()// قراءة عمود JSON كـ VariantType
Dataset<Row> df = spark.sql("SELECT id, data FROM clickhouse.default.json_table");
// الوصول إلى بيانات Variant
df.show();
// تحويل Variant إلى سلسلة JSON للفحص
import static org.apache.spark.sql.functions.*;
df.select(
col("id"),
to_json(col("data")).as("data_json")
).show();كتابة بيانات من النوع VariantType
يمكنك كتابة بيانات VariantType إلى ClickHouse باستخدام نوعَي الأعمدة JSON أو Variant:
import org.apache.spark.sql.functions._
// إنشاء DataFrame يحتوي على بيانات JSON
val jsonData = Seq(
(1, """{"name": "Alice", "age": 30}"""),
(2, """{"name": "Bob", "age": 25}"""),
(3, """{"name": "Charlie", "city": "NYC"}""")
).toDF("id", "json_string")
// تحليل سلاسل JSON إلى VariantType
val variantDF = jsonData.select(
col("id"),
parse_json(col("json_string")).as("data")
)
// الكتابة إلى ClickHouse باستخدام نوع JSON (كائنات JSON فقط)
variantDF.writeTo("clickhouse.default.user_data").create()
// أو حدِّد Variant بعدة أنواع
spark.sql("""
CREATE TABLE clickhouse.default.mixed_data (
id INT,
data VARIANT
) USING clickhouse
TBLPROPERTIES (
'clickhouse.column.data.variant_types' = 'String, Int64, Bool, JSON',
'engine' = 'MergeTree()',
'order_by' = 'id'
)
""")from pyspark.sql.functions import parse_json
# إنشاء DataFrame يحتوي على بيانات JSON
json_data = [
(1, '{"name": "Alice", "age": 30}'),
(2, '{"name": "Bob", "age": 25}'),
(3, '{"name": "Charlie", "city": "NYC"}')
]
df = spark.createDataFrame(json_data, ["id", "json_string"])
# تحليل سلاسل JSON إلى VariantType
variant_df = df.select(
"id",
parse_json("json_string").alias("data")
)
# الكتابة إلى ClickHouse باستخدام نوع JSON
variant_df.writeTo("clickhouse.default.user_data").create()
# أو حدِّد Variant بعدة أنواع
spark.sql("""
CREATE TABLE clickhouse.default.mixed_data (
id INT,
data VARIANT
) USING clickhouse
TBLPROPERTIES (
'clickhouse.column.data.variant_types' = 'String, Int64, Bool, JSON',
'engine' = 'MergeTree()',
'order_by' = 'id'
)
""")import static org.apache.spark.sql.functions.*;
// إنشاء DataFrame يحتوي على بيانات JSON
List<Row> jsonData = Arrays.asList(
RowFactory.create(1, "{\"name\": \"Alice\", \"age\": 30}"),
RowFactory.create(2, "{\"name\": \"Bob\", \"age\": 25}"),
RowFactory.create(3, "{\"name\": \"Charlie\", \"city\": \"NYC\"}")
);
StructType schema = new StructType(new StructField[]{
DataTypes.createStructField("id", DataTypes.IntegerType, false),
DataTypes.createStructField("json_string", DataTypes.StringType, false)
});
Dataset<Row> jsonDF = spark.createDataFrame(jsonData, schema);
// تحليل سلاسل JSON إلى VariantType
Dataset<Row> variantDF = jsonDF.select(
col("id"),
parse_json(col("json_string")).as("data")
);
// الكتابة إلى ClickHouse باستخدام نوع JSON (كائنات JSON فقط)
variantDF.writeTo("clickhouse.default.user_data").create();
// أو حدِّد Variant بعدة أنواع
spark.sql("CREATE TABLE clickhouse.default.mixed_data (" +
"id INT, " +
"data VARIANT" +
") USING clickhouse " +
"TBLPROPERTIES (" +
"'clickhouse.column.data.variant_types' = 'String, Int64, Bool, JSON', " +
"'engine' = 'MergeTree()', " +
"'order_by' = 'id'" +
")");إنشاء جداول من نوع VariantType باستخدام Spark SQL
يمكنك إنشاء جداول من نوع VariantType باستخدام عبارات DDL في Spark SQL:
-- Create table with JSON type (default)
CREATE TABLE clickhouse.default.json_table (
id INT,
data VARIANT
) USING clickhouse
TBLPROPERTIES (
'engine' = 'MergeTree()',
'order_by' = 'id'
)-- Create table with Variant type supporting multiple types
CREATE TABLE clickhouse.default.flexible_data (
id INT,
data VARIANT
) USING clickhouse
TBLPROPERTIES (
'clickhouse.column.data.variant_types' = 'String, Int64, Float64, Bool, Array(String), JSON',
'engine' = 'MergeTree()',
'order_by' = 'id'
)إعداد أنواع Variant
عند إنشاء جداول تضم أعمدة من نوع VariantType، يمكنك تحديد أنواع ClickHouse التي تريد استخدامها:
نوع JSON (الافتراضي)
إذا لم تُحدَّد الخاصية variant_types، فسيستخدم العمود نوع JSON الافتراضي في ClickHouse، والذي لا يقبل إلا كائنات JSON:
CREATE TABLE clickhouse.default.json_table (
id INT,
data VARIANT
) USING clickhouse
TBLPROPERTIES (
'engine' = 'MergeTree()',
'order_by' = 'id'
)ينشئ هذا استعلام ClickHouse التالي:
CREATE TABLE json_table (id Int32, data JSON) ENGINE = MergeTree() ORDER BY idنوع Variant مع أنواع متعددة
لدعم الأنواع الأولية والمصفوفات وكائنات JSON، حدِّد الأنواع في الخاصية variant_types:
CREATE TABLE clickhouse.default.flexible_data (
id INT,
data VARIANT
) USING clickhouse
TBLPROPERTIES (
'clickhouse.column.data.variant_types' = 'String, Int64, Float64, Bool, Array(String), JSON',
'engine' = 'MergeTree()',
'order_by' = 'id'
)ينشئ هذا استعلام ClickHouse التالي:
CREATE TABLE flexible_data (
id Int32,
data Variant(String, Int64, Float64, Bool, Array(String), JSON)
) ENGINE = MergeTree() ORDER BY idأنواع Variant المدعومة
يمكن استخدام أنواع ClickHouse التالية في Variant():
- الأنواع الأساسية:
StringوInt8وInt16وInt32وInt64وUInt8وUInt16وUInt32وUInt64وFloat32وFloat64وBool - المصفوفات:
Array(T)، حيث يكون T أي نوع مدعوم، بما في ذلك المصفوفات المتداخلة - JSON:
JSONلتخزين كائنات JSON
إعداد تنسيق القراءة
بشكل افتراضي، تُقرأ أعمدة JSON وVariant على أنها VariantType. يمكنك تجاوز هذا السلوك وقراءتها كسلاسل نصية:
// قراءة JSON/Variant كسلاسل نصية بدلًا من VariantType
spark.conf.set("spark.clickhouse.read.jsonAs", "string")
val df = spark.sql("SELECT id, data FROM clickhouse.default.json_table")
// سيكون العمود data من النوع StringType ويحتوي على سلاسل JSON# قراءة JSON/Variant كسلاسل نصية بدلًا من VariantType
spark.conf.set("spark.clickhouse.read.jsonAs", "string")
df = spark.sql("SELECT id, data FROM clickhouse.default.json_table")
# سيكون العمود data من النوع StringType ويحتوي على سلاسل JSON// قراءة JSON/Variant كسلاسل نصية بدلًا من VariantType
spark.conf().set("spark.clickhouse.read.jsonAs", "string");
Dataset<Row> df = spark.sql("SELECT id, data FROM clickhouse.default.json_table");
// سيكون العمود data من النوع StringType ويحتوي على سلاسل JSONدعم صيغ الكتابة
يختلف دعم الكتابة لـ VariantType باختلاف الصيغة:
| التنسيق | الدعم | ملاحظات |
|---|---|---|
| JSON | ✅ كامل | يدعم النوعين JSON وVariant كليهما. يُوصى به لبيانات VariantType |
| Arrow | ⚠️ جزئي | يدعم الكتابة إلى النوع JSON في ClickHouse، لكنه لا يدعم النوع Variant في ClickHouse. الدعم الكامل بانتظار حل المشكلة https://github.com/ClickHouse/ClickHouse/issues/92752 |
اضبط صيغة الكتابة:
spark.conf.set("spark.clickhouse.write.format", "json") // Recommended for Variant typesأفضل الممارسات
- استخدم نوع JSON للبيانات التي تحتوي على JSON فقط: إذا كنت تخزّن كائنات JSON فقط، فاستخدم نوع JSON الافتراضي (من دون الخاصية
variant_types) - حدّد الأنواع صراحةً: عند استخدام
Variant()، اذكر صراحةً جميع الأنواع التي تنوي تخزينها - فعّل الميزات التجريبية: تأكد من تفعيل
allow_experimental_json_type = 1في ClickHouse - استخدم تنسيق JSON لعمليات الكتابة: يُوصى باستخدام تنسيق JSON لبيانات VariantType لتحقيق توافق أفضل
- ضع أنماط الاستعلام في الاعتبار: تدعم أنواع JSON/Variant استعلامات مسار JSON في ClickHouse لتصفية أكثر كفاءة
- تلميحات الأعمدة لتحسين الأداء: عند استخدام حقول JSON في ClickHouse، يؤدي إضافة تلميحات الأعمدة إلى تحسين أداء الاستعلامات. لا تتوفر حاليًا إمكانية إضافة تلميحات الأعمدة عبر Spark. راجع GitHub issue #497 لمتابعة هذه الميزة.
مثال: سير عمل متكامل
import org.apache.spark.sql.functions._
// Enable experimental JSON type in ClickHouse
spark.sql("SET allow_experimental_json_type = 1")
// Create table with Variant column
spark.sql("""
CREATE TABLE clickhouse.default.events (
event_id BIGINT,
event_time TIMESTAMP,
event_data VARIANT
) USING clickhouse
TBLPROPERTIES (
'clickhouse.column.event_data.variant_types' = 'String, Int64, Bool, JSON',
'engine' = 'MergeTree()',
'order_by' = 'event_time'
)
""")
// Prepare data with mixed types
val events = Seq(
(1L, "2024-01-01 10:00:00", """{"action": "login", "user_id": 123}"""),
(2L, "2024-01-01 10:05:00", """{"action": "purchase", "amount": 99.99}"""),
(3L, "2024-01-01 10:10:00", """{"action": "logout", "duration": 600}""")
).toDF("event_id", "event_time", "json_data")
// Convert to VariantType and write
val variantEvents = events.select(
col("event_id"),
to_timestamp(col("event_time")).as("event_time"),
parse_json(col("json_data")).as("event_data")
)
variantEvents.writeTo("clickhouse.default.events").append()
// Read and query
val result = spark.sql("""
SELECT event_id, event_time, event_data
FROM clickhouse.default.events
WHERE event_time >= '2024-01-01'
ORDER BY event_time
""")
result.show(false)from pyspark.sql.functions import parse_json, to_timestamp
# Enable experimental JSON type in ClickHouse
spark.sql("SET allow_experimental_json_type = 1")
# Create table with Variant column
spark.sql("""
CREATE TABLE clickhouse.default.events (
event_id BIGINT,
event_time TIMESTAMP,
event_data VARIANT
) USING clickhouse
TBLPROPERTIES (
'clickhouse.column.event_data.variant_types' = 'String, Int64, Bool, JSON',
'engine' = 'MergeTree()',
'order_by' = 'event_time'
)
""")
# Prepare data with mixed types
events = [
(1, "2024-01-01 10:00:00", '{"action": "login", "user_id": 123}'),
(2, "2024-01-01 10:05:00", '{"action": "purchase", "amount": 99.99}'),
(3, "2024-01-01 10:10:00", '{"action": "logout", "duration": 600}')
]
df = spark.createDataFrame(events, ["event_id", "event_time", "json_data"])
# Convert to VariantType and write
variant_events = df.select(
"event_id",
to_timestamp("event_time").alias("event_time"),
parse_json("json_data").alias("event_data")
)
variant_events.writeTo("clickhouse.default.events").append()
# Read and query
result = spark.sql("""
SELECT event_id, event_time, event_data
FROM clickhouse.default.events
WHERE event_time >= '2024-01-01'
ORDER BY event_time
""")
result.show(truncate=False)import static org.apache.spark.sql.functions.*;
// Enable experimental JSON type in ClickHouse
spark.sql("SET allow_experimental_json_type = 1");
// Create table with Variant column
spark.sql("CREATE TABLE clickhouse.default.events (" +
"event_id BIGINT, " +
"event_time TIMESTAMP, " +
"event_data VARIANT" +
") USING clickhouse " +
"TBLPROPERTIES (" +
"'clickhouse.column.event_data.variant_types' = 'String, Int64, Bool, JSON', " +
"'engine' = 'MergeTree()', " +
"'order_by' = 'event_time'" +
")");
// Prepare data with mixed types
List<Row> events = Arrays.asList(
RowFactory.create(1L, "2024-01-01 10:00:00", "{\"action\": \"login\", \"user_id\": 123}"),
RowFactory.create(2L, "2024-01-01 10:05:00", "{\"action\": \"purchase\", \"amount\": 99.99}"),
RowFactory.create(3L, "2024-01-01 10:10:00", "{\"action\": \"logout\", \"duration\": 600}")
);
StructType eventSchema = new StructType(new StructField[]{
DataTypes.createStructField("event_id", DataTypes.LongType, false),
DataTypes.createStructField("event_time", DataTypes.StringType, false),
DataTypes.createStructField("json_data", DataTypes.StringType, false)
});
Dataset<Row> eventsDF = spark.createDataFrame(events, eventSchema);
// Convert to VariantType and write
Dataset<Row> variantEvents = eventsDF.select(
col("event_id"),
to_timestamp(col("event_time")).as("event_time"),
parse_json(col("json_data")).as("event_data")
);
variantEvents.writeTo("clickhouse.default.events").append();
// Read and query
Dataset<Row> result = spark.sql("SELECT event_id, event_time, event_data " +
"FROM clickhouse.default.events " +
"WHERE event_time >= '2024-01-01' " +
"ORDER BY event_time");
result.show(false);الإعدادات
فيما يلي الإعدادات القابلة للتعديل المتاحة في الموصّل.
| المفتاح | القيمة الافتراضية | الوصف | منذ الإصدار |
|---|---|---|---|
| spark.clickhouse.ignoreUnsupportedTransform | true | يدعم ClickHouse استخدام تعبيرات معقدة كمفاتيح تجزئة أو قيم partition، مثل cityHash64(col_1, col_2)، لكن Spark لا يدعمها حاليًا. إذا كانت القيمة true، فتجاهل التعبيرات غير المدعومة وسجّل تحذيرًا، وإلا فأخفِق فورًا مع استثناء. تحذير: عند ضبط spark.clickhouse.write.distributed.convertLocal=true، قد يؤدي تجاهل مفاتيح التجزئة غير المدعومة إلى تلف البيانات. يتحقق الموصّل من ذلك ويُصدر خطأً افتراضيًا. وللسماح بذلك، اضبط spark.clickhouse.write.distributed.convertLocal.allowUnsupportedSharding=true صراحةً. |
0.4.0 |
| spark.clickhouse.read.compression.codec | lz4 | خوارزمية الضغط المستخدمة لفك ضغط البيانات عند القراءة. خوارزميات الضغط المدعومة: none وlz4. | 0.5.0 |
| spark.clickhouse.read.distributed.convertLocal | true | عند القراءة من جدول Distributed، يُقرأ الجدول المحلي بدلًا منه. إذا كانت القيمة true، يتم تجاهل spark.clickhouse.read.distributed.useClusterNodes. |
0.1.0 |
| spark.clickhouse.read.fixedStringAs | binary | اقرأ نوع FixedString في ClickHouse كنوع بيانات Spark المحدد. الأنواع المدعومة: binary, string | 0.8.0 |
| spark.clickhouse.read.format | json | تنسيق التسلسل المستخدم للقراءة. التنسيقات المدعومة: json وbinary | 0.6.0 |
| spark.clickhouse.read.runtimeFilter.enabled | false | تمكين مرشح وقت التشغيل لعمليات القراءة. | 0.8.0 |
| spark.clickhouse.read.splitByPartitionId | true | إذا كانت القيمة true، فأنشئ عامل تصفية قسم الإدخال باستخدام العمود الافتراضي _partition_id بدلًا من قيمة القسم. توجد مشكلات معروفة في بناء شروط SQL استنادًا إلى قيمة القسم. تتطلب هذه الميزة ClickHouse Server v21.6+ |
0.4.0 |
| spark.clickhouse.useNullableQuerySchema | false | إذا كانت القيمة true، فاعتبر جميع حقول مخطط الاستعلام قابلة لـ NULL عند تنفيذ CREATE/REPLACE TABLE ... AS SELECT ... أثناء إنشاء الجدول. لاحظ أن هذا الإعداد يتطلب SPARK-43390 (المتوفر في Spark 3.5)، وبدون هذا التصحيح، فإنه يتصرف دائمًا كما لو كانت قيمته true. |
0.8.0 |
| spark.clickhouse.write.batchSize | 10000 | عدد السجلات لكل دفعة عند الكتابة إلى ClickHouse. | 0.1.0 |
| spark.clickhouse.write.compression.codec | lz4 | خوارزمية الضغط المستخدمة عند كتابة البيانات. خوارزميات الضغط المدعومة: none وlz4. | 0.3.0 |
| spark.clickhouse.write.distributed.convertLocal | false | عند الكتابة إلى جدول Distributed، تتم الكتابة إلى الجدول المحلي بدلًا منه. إذا كانت القيمة true، يتم تجاهل spark.clickhouse.write.distributed.useClusterNodes. يتجاوز هذا آلية التوجيه الأصلية في ClickHouse، ما يتطلب من Spark تقييم مفتاح التقسيم. عند استخدام تعبيرات تقسيم غير مدعومة، اضبط spark.clickhouse.ignoreUnsupportedTransform على false لمنع حدوث أخطاء صامتة في توزيع البيانات. |
0.1.0 |
| spark.clickhouse.write.distributed.convertLocal.allowUnsupportedSharding | false | اسمح بالكتابة إلى جداول Distributed مع convertLocal=true وignoreUnsupportedTransform=true عندما يكون مفتاح التشظية غير مدعوم. هذا خيار خطير وقد يؤدي إلى تلف البيانات بسبب التشظية غير الصحيحة. عند ضبطه على true، يجب التأكد من أن بياناتك مرتبة/مُشظّاة بشكل صحيح قبل الكتابة، لأن Spark لا يمكنه تقييم تعبير التشظية غير المدعوم. لا تضبطه على true إلا إذا كنت تدرك المخاطر وقد تحققت من توزيع بياناتك. افتراضيًا، سيؤدي هذا الجمع إلى ظهور خطأ لمنع تلف البيانات الصامت. |
0.10.0 |
| spark.clickhouse.write.distributed.useClusterNodes | true | الكتابة إلى كل عُقد العنقود عند الكتابة إلى جدول Distributed. | 0.1.0 |
| spark.clickhouse.write.format | arrow | صيغة التسلسل المستخدمة للكتابة. التنسيقات المدعومة: json وarrow | 0.4.0 |
| spark.clickhouse.write.localSortByKey | true | إذا كانت القيمة true، فسيُجرى فرز محلي وفقًا لمفاتيح الفرز قبل الكتابة. |
0.3.0 |
| spark.clickhouse.write.localSortByPartition | القيمة الخاصة بـ spark.clickhouse.write.repartitionByPartition | إذا كانت القيمة true، فأجرِ فرزًا محليًا حسب التقسيم قبل الكتابة. وإذا لم تُعيَّن، فستكون مساويةً لـ spark.clickhouse.write.repartitionByPartition. |
0.3.0 |
| spark.clickhouse.write.maxRetry | 3 | الحد الأقصى لعدد مرات إعادة محاولة عملية كتابة دفعة واحدة إذا فشلت بسبب رموز خطأ قابلة لإعادة المحاولة. | 0.1.0 |
| spark.clickhouse.write.repartitionByPartition | true | ما إذا كان ينبغي إعادة تقسيم البيانات بحسب مفاتيح تقسيم ClickHouse لتتوافق مع توزيعات جدول ClickHouse قبل الكتابة. | 0.3.0 |
| spark.clickhouse.write.repartitionNum | 0 | يلزم إعادة تقسيم البيانات بما يتوافق مع توزيع جدول ClickHouse قبل الكتابة. استخدم هذا الإعداد لتحديد عدد أقسام إعادة التقسيم؛ وأي قيمة أقل من 1 تعني عدم الحاجة إلى ذلك. | 0.1.0 |
| spark.clickhouse.write.repartitionStrictly | false | إذا كانت true، فسيُوزّع Spark السجلات الواردة بدقة عبر الأقسام لتلبية التوزيع المطلوب قبل تمريرها إلى جدول مصدر البيانات عند الكتابة. وإلا، فقد يطبّق Spark بعض التحسينات لتسريع الاستعلام، لكن على حساب متطلب التوزيع. ملاحظة: يتطلب هذا الإعداد SPARK-37523 (المتوفرة في Spark 3.4)، وبدون هذا التصحيح سيعمل دائمًا كما لو كانت قيمته true. |
0.3.0 |
| spark.clickhouse.write.retryInterval | 10s | الفاصل الزمني، بالثواني، بين محاولات إعادة الكتابة. | 0.1.0 |
| spark.clickhouse.write.retryableErrorCodes | 241 | رموز الخطأ القابلة لإعادة المحاولة التي يعيدها خادم ClickHouse عند فشل عملية الكتابة. | 0.1.0 |
أنواع البيانات المدعومة
يوضح هذا القسم كيفية توافق أنواع البيانات بين Spark وClickHouse. وتوفر الجداول أدناه مراجع سريعة لتحويل أنواع البيانات عند قراءتها من ClickHouse إلى Spark، وعند إدراجها من Spark في ClickHouse.
قراءة البيانات من ClickHouse إلى Spark
| نوع بيانات ClickHouse | نوع بيانات Spark | مدعوم | بدائي | ملاحظات |
|---|---|---|---|---|
Nothing |
NullType |
✅ | نعم | |
Bool |
BooleanType |
✅ | نعم | |
UInt8, Int16 |
ShortType |
✅ | نعم | |
Int8 |
ByteType |
✅ | نعم | |
UInt16,Int32 |
IntegerType |
✅ | نعم | |
UInt32,Int64, UInt64 |
LongType |
✅ | نعم | |
Int128,UInt128, Int256, UInt256 |
DecimalType(38, 0) |
✅ | نعم | |
Float32 |
FloatType |
✅ | نعم | |
Float64 |
DoubleType |
✅ | نعم | |
String, UUID, Enum8, Enum16, IPv4, IPv6 |
StringType |
✅ | نعم | |
FixedString |
BinaryType, StringType |
✅ | نعم | يتحكم فيه الإعداد READ_FIXED_STRING_AS |
Decimal |
DecimalType |
✅ | نعم | حتى دقة ومقياس Decimal128 |
Decimal32 |
DecimalType(9, scale) |
✅ | نعم | |
Decimal64 |
DecimalType(18, scale) |
✅ | نعم | |
Decimal128 |
DecimalType(38, scale) |
✅ | نعم | |
Date, Date32 |
DateType |
✅ | نعم | |
DateTime, DateTime32, DateTime64 |
TimestampType |
✅ | نعم | |
Array |
ArrayType |
✅ | لا | يُحوَّل أيضًا نوع عنصر الـArray |
Map |
MapType |
✅ | لا | المفاتيح مقتصرة على StringType |
IntervalYear |
YearMonthIntervalType(Year) |
✅ | نعم | |
IntervalMonth |
YearMonthIntervalType(Month) |
✅ | نعم | |
IntervalDay, IntervalHour, IntervalMinute, IntervalSecond |
DayTimeIntervalType |
✅ | لا | يُستخدم نوع interval المحدد |
JSON, Variant |
VariantType |
✅ | لا | يتطلب Spark 4.0+ وClickHouse 25.3+. ويمكن قراءته كـ StringType باستخدام spark.clickhouse.read.jsonAs=string |
Object |
❌ | |||
Nested |
❌ | |||
Tuple |
StructType |
✅ | لا | يدعم كلاً من Tuple المسمّاة وغير المسمّاة. تُطابَق الـTuple المسمّاة مع حقول struct حسب الاسم، بينما تستخدم الـTuple غير المسمّاة _1 و_2 وما إلى ذلك. كما يدعم structs المتداخلة والحقول القابلة للقيمة NULL |
Point |
❌ | |||
Polygon |
❌ | |||
MultiPolygon |
❌ | |||
Ring |
❌ | |||
IntervalQuarter |
❌ | |||
IntervalWeek |
❌ | |||
Decimal256 |
❌ | |||
AggregateFunction |
❌ | |||
SimpleAggregateFunction |
❌ |
إدراج البيانات من Spark إلى ClickHouse
| نوع بيانات Spark | نوع بيانات ClickHouse | مدعوم | بدائي | ملاحظات |
|---|---|---|---|---|
BooleanType |
Bool |
✅ | نعم | يُعيَّن إلى النوع Bool (وليس UInt8) بدءًا من الإصدار 0.9.0 |
ByteType |
Int8 |
✅ | نعم | |
ShortType |
Int16 |
✅ | نعم | |
IntegerType |
Int32 |
✅ | نعم | |
LongType |
Int64 |
✅ | نعم | |
FloatType |
Float32 |
✅ | نعم | |
DoubleType |
Float64 |
✅ | نعم | |
StringType |
String |
✅ | نعم | |
VarcharType |
String |
✅ | نعم | |
CharType |
String |
✅ | نعم | |
DecimalType |
Decimal(p, s) |
✅ | نعم | الدقة والمقياس حتى Decimal128 |
DateType |
Date |
✅ | نعم | |
TimestampType |
DateTime |
✅ | نعم | |
ArrayType (list, tuple, or array) |
Array |
✅ | لا | يُحوَّل أيضًا نوع عنصر المصفوفة |
MapType |
Map |
✅ | لا | المفاتيح مقتصرة على StringType |
StructType |
Tuple |
✅ | لا | يُحوَّل إلى Tuple مُسمّى باستخدام أسماء الحقول. |
VariantType |
JSON or Variant |
✅ | لا | يتطلب Spark 4.0+ وClickHouse 25.3+. النوع الافتراضي هو JSON. استخدم الخاصية clickhouse.column.<name>.variant_types لتحديد Variant مع أنواع متعددة. |
Object |
❌ | |||
Nested |
❌ |
المساهمة والدعم
إذا كنت ترغب في المساهمة في المشروع أو الإبلاغ عن أي مشكلات، فنحن نرحب بمشاركتك! زُر مستودع GitHub الخاص بنا لفتح مشكلة، أو اقتراح تحسينات، أو إرسال طلب سحب. نرحب بالمساهمات! يُرجى مراجعة إرشادات المساهمة في المستودع قبل البدء. شكرًا لمساعدتك في تحسين ClickHouse Spark connector الخاص بنا!