Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
İşlem hatlarında kullanarak
Önemli
Bu özellik Genel Önizleme aşamasındadır.
Lakeflow işlem hatlarında, from_json SQL işlevi siz açık bir şema sağlamadan JSON bloblarının şemasını otomatik olarak çıkarıp geliştirebilir.
from_json işlem hatlarında nasıl çalışır
from_json SQL işlevi bir JSON dize sütununu ayrıştırıp bir yapı değeri döndürür. İşlem hattının dışında kullanıldığında, schema argümanını kullanarak döndürülen değerin şemasını açıkça sağlamanız gerekir. İşlem hattında kullanıldığında, döndürülen değerin şemasını otomatik olarak yöneten şema çıkarımı ve evrimini etkinleştirebilirsiniz. Bu özellik, hem ilk kurulumu (özellikle şema bilinmiyorsa) hem de şema sık sık değiştiğinde devam eden işlemleri basitleştirir.
Otomatik Yükleyici, Kafka veya Kinesis gibi akış veri kaynaklarından rastgele JSON bloblarını işler.
Özellikle, bir işlem hattında kullanıldığında, SQL işlevi için şema çıkarımı ve evrimi from_json şunları yapabilir:
- Gelen JSON kayıtlarındaki yeni alanları algılama (iç içe JSON nesneleri dahil)
- Alan türlerini belirleyin ve uygun Spark veri tipleriyle eşleştirin
- Şemayı yeni alanlara uyum sağlamak için otomatik olarak geliştirin
- Geçerli şemaya uymayan verileri otomatik olarak işleme
Söz dizimi: Şemayı otomatik olarak çıkar ve geliştir
Şema çıkarımını from_json kullanarak bir işlem hattında etkinleştirmek için, şemayı NULL olarak ayarlayın ve schemaLocationKey seçeneğini belirtin. Bu, şemayı çıkarıp izlemesini sağlar.
SQL
from_json(jsonStr, NULL, map("schemaLocationKey", "<uniqueKey>” [, otherOptions]))
Piton
from_json(jsonStr, None, {"schemaLocationKey": "<uniqueKey>”[, otherOptions]})
Sorguda birden çok from_json ifade olabilir, ancak her ifadenin benzersiz schemaLocationKeybir ifadesi olmalıdır. ayrıca schemaLocationKey işlem hattı başına benzersiz olmalıdır.
SQL
SELECT
value,
from_json(value, NULL, map('schemaLocationKey', 'keyX')) parsedX,
from_json(value, NULL, map('schemaLocationKey', 'keyY')) parsedY,
FROM STREAM READ_FILES('/databricks-datasets/nyctaxi/sample/json/', format => 'text')
Piton
(spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "text")
.load("/databricks-datasets/nyctaxi/sample/json/")
.select(
col("value"),
from_json(col("value"), None, {"schemaLocationKey": "keyX"}).alias("parsedX"),
from_json(col("value"), None, {"schemaLocationKey": "keyY"}).alias("parsedY"))
)
Sözdizimi: Sabit şema
Bunun yerine belirli bir şemayı zorunlu kılmak istiyorsanız, JSON dizesini bu şemayı kullanarak ayrıştırmak için aşağıdaki from_json söz dizimini kullanabilirsiniz:
from_json(jsonStr, schema, [, options])
Bu söz dizimi, işlem hatları dahil olmak üzere herhangi bir Azure Databricks ortamında kullanılabilir. burada
Şema Çıkarımı
from_json JSON veri sütunlarının ilk toplu işleminden şemayı çıkarır ve dahili olarak schemaLocationKey'e göre dizine alır (gerekli).
JSON dizesi tek bir nesneyse (örneğin, {"id": 123, "name": "John"}), from_json STRUCT türünde bir şema çıkartır ve alan listesine bir rescuedDataColumn ekler.
STRUCT<id LONG, name STRING, _rescued_data STRING>
Ancak, JSON dizesi bir üst düzey dizi (örneğin ["id": 123, "name": "John"]) içeriyorsa, from_json ARRAY'ı STRUCT içinde sarar. Bu yaklaşım, çıkarsanan şemayla uyumlu olmayan verilerin çıkarıldığını etkinleştirir. Dizi değerlerini ayrı satırlara ayırarak ileriki adımlarda kullanma seçeneğiniz vardır.
STRUCT<value ARRAY<id LONG, name STRING>, _rescued_data STRING>
Şema ipuçlarını kullanarak şema çıkarımı geçersiz kılma
İsteğe bağlı olarak, schemaHints sağlayarak from_json'in bir sütunun türünü nasıl çıkaracağını etkileyebilirsiniz. Bu, bir sütunun belirli bir veri türünde olduğunu bildiğinizde veya daha genel bir veri türü (örneğin, tamsayı yerine çift) seçmek istediğinizde yararlıdır. SQL şema belirtimi söz dizimini kullanarak sütun veri türleri için rastgele sayıda ipucu sağlayabilirsiniz. Şema ipuçlarının semantiği, Otomatik Yükleyici şema ipuçlarıyla aynıdır. Örneğin:
SELECT
-- The JSON `{"a": 1}` will treat `a` as a BIGINT
from_json(data, NULL, map('schemaLocationKey', 'w', 'schemaHints', '')),
-- The JSON `{"a": 1}` will treat `a` as a STRING
from_json(data, NULL, map('schemaLocationKey', 'x', 'schemaHints', 'a STRING')),
-- The JSON `{"a": {"b": 1}}` will treat `a` as a MAP<STRING, BIGINT>
from_json(data, NULL, map('schemaLocationKey', 'y', 'schemaHints', 'a MAP<STRING, BIGINT'>)),
-- The JSON `{"a": {"b": 1}}` will treat `a` as a STRING
from_json(data, NULL, map('schemaLocationKey', 'z', 'schemaHints', 'a STRING')),
FROM STREAM READ_FILES(...)
JSON dizesi üst düzey bir ARRAY içeriyorsa, bir STRUCT içinde sarmalanır. Bu durumlarda, sarmalanan STRUCT yerine ARRAY şemasına şema ipuçları uygulanır. Örneğin, aşağıdaki gibi bir üst düzey diziye sahip bir JSON dizesi düşünün:
[{"id": 123, "name": "John"}]
Çıkarsanan ARRAY şeması, bir STRUCT içinde sarmalanır.
STRUCT<value ARRAY<id LONG, name STRING>, _rescued_data STRING>
veri türünü iddeğiştirmek için şema ipucunu STRING olarak element.id belirtin. DOUBLE türünde yeni bir sütun eklemek için element.new_col DOUBLE olarak tanımlayın. Bu ipuçları nedeniyle en üst düzey JSON dizisinin şeması şöyle olur:
struct<value array<id STRING, name STRING, new_col DOUBLE>, _rescued_data STRING>
Şemayı schemaEvolutionMode kullanarak evrimleştirin
from_json verilerinizi işlerken yeni sütunların eklenmesini algılar. Yeni bir alan algıladığında from_json , yeni sütunları şemanın sonuna birleştirerek çıkarılan şemayı en son şemayla güncelleştirir. Mevcut sütunların veri türleri değişmeden kalır. Şema güncelleştirmesinin ardından işlem hattı güncelleştirilmiş şemayla otomatik olarak yeniden başlatılır.
from_json , isteğe bağlı schemaEvolutionMode ayarı kullanarak ayarladığınız şema evrimi için aşağıdaki modları destekler. Bu modlar Otomatik Yükleyici ile tutarlıdır.
schemaEvolutionMode |
Yeni sütun okuma davranışı |
|---|---|
addNewColumns (varsayılan) |
Akış başarısız oluyor. Şemaya yeni sütunlar eklenir. Mevcut sütunlar veri türlerini geliştirmez. |
rescue |
Şema hiçbir zaman geliştirilmez ve şema değişiklikleri nedeniyle akış başarısız olmaz. Tüm yeni sütunlar kurtarılan veri sütununa kaydedilir. |
failOnNewColumns |
Akış başarısız oluyor. Akış güncelleştirilmediği veya sorunlu veriler kaldırılmadığı sürece schemaHints yeniden başlatılmaz. |
none |
Şemayı evrimleştirmez, yeni sütunlar yoksayılır ve rescuedDataColumn seçeneği ayarlanmadıkça veriler kurtarılmaz. Şema değişiklikleri nedeniyle akış başarısız olmaz. |
Örneğin:
SELECT
-- If a new column appears, the pipeline will automatically add it to the schema:
from_json(a, NULL, map('schemaLocationKey', 'w', 'schemaEvolutionMode', 'addNewColumns')),
-- If a new column appears, the pipeline will add it to the rescued data column:
from_json(b, NULL, map('schemaLocationKey', 'x', 'schemaEvolutionMode', 'rescue')),
-- If a new column appears, the pipeline will ignore it:
from_json(c, NULL, map('schemaLocationKey', 'y', 'schemaEvolutionMode', 'none')),
-- If a new column appears, the pipeline will fail:
from_json(d, NULL, map('schemaLocationKey', 'z', 'schemaEvolutionMode', 'failOnNewColumns')),
FROM STREAM READ_FILES(...)
Kurtarılan veri sütunu
Kurtarılan veri sütunu şemanıza otomatik olarak olarak _rescued_dataeklenir. seçeneğini ayarlayarak rescuedDataColumn sütunu yeniden adlandırabilirsiniz. Örneğin:
from_json(jsonStr, None, {"schemaLocationKey": "keyX", "rescuedDataColumn": "my_rescued_data"})
Kurtarılan veri sütununu kullanmayı seçtiğinizde, çıkarılmış şemayla eşleşmeyen tüm sütunlar bırakılmak yerine kurtarılır. Bu durum, bir veri türü uyuşmazlığından, şemadaki eksik bir sütundan veya sütun adı büyük/küçük harf farkından kaynaklanıyor olabilir.
Bozuk kayıtları işleme
Hatalı biçimlendirilmiş ve ayrıştırılamayan kayıtları depolamak için, aşağıdaki örnekte olduğu gibi şema ipuçlarını ayarlayarak bir _corrupt_record sütun ekleyin:
CREATE STREAMING TABLE bronze AS
SELECT
from_json(value, NULL,
map('schemaLocationKey', 'nycTaxi',
'schemaHints', '_corrupt_record STRING',
'columnNameOfCorruptRecord', '_corrupt_record')) jsonCol
FROM STREAM READ_FILES('/databricks-datasets/nyctaxi/sample/json/', format => 'text')
Bozuk kayıt sütununu yeniden adlandırmak columnNameOfCorruptRecord için seçeneğini ayarlayın.
JSON ayrıştırıcısı bozuk kayıtları işlemek için üç modu destekler:
| Mode | Description |
|---|---|
PERMISSIVE |
Bozuk kayıtlar için, hatalı biçimlendirilmiş dizeyi columnNameOfCorruptRecord ile yapılandırılan bir alana yerleştirir ve hatalı biçimlendirilmiş alanları null olarak ayarlar. Bozuk kayıtları tutmak için, kullanıcı tanımlı şemada adlı columnNameOfCorruptRecord bir dize türü alanı ayarlayabilirsiniz. Bir şemada alan yoksa, ayrıştırma sırasında bozuk kayıtlar bırakılır. Bir şema çıkarıldığında ayrıştırıcı, çıkış şemasına örtük olarak bir columnNameOfCorruptRecord alan ekler. |
DROPMALFORMED |
Bozuk kayıtları yok sayar.DROPMALFORMED ve rescuedDataColumn modunu kullandığınızda, veri türü uyuşmazlıkları kayıtların bırakılmasına neden olmaz. Eksik veya hatalı biçimlendirilmiş JSON gibi yalnızca bozuk kayıtlar bırakılır. |
FAILFAST |
Ayrıştırıcı bozuk kayıtları karşıladığında bir özel durum oluşturur. ile FAILFASTmodu kullandığınızdarescuedDataColumn, veri türü uyuşmazlıkları hata oluşturmaz. Yalnızca bozuk kayıtlar eksik veya hatalı biçimlendirilmiş JSON gibi hatalar oluşturur. |
from_json çıktısında bir alana bakın
from_json işlem hattı yürütmesi sırasında şemayı çıkarsar. Aşağı akış sorgusu, işlev en az bir kez başarıyla yürütülmeden önce from_json bir from_json alana başvuruyorsa, alan çözümlenmez ve sorgu atlanır. Aşağıdaki örnekte, bronz sorgudaki işlev yürütülene ve şema çıkarılana kadar from_json gümüş tablo sorgusu analizi atlanır.
CREATE STREAMING TABLE bronze AS
SELECT
from_json(value, NULL, map('schemaLocationKey', 'nycTaxi')) jsonCol
FROM STREAM READ_FILES('/databricks-datasets/nyctaxi/sample/json/', format => 'text')
CREATE STREAMING TABLE silver AS
SELECT jsonCol.VendorID, jsonCol.total_amount
FROM bronze
İşlev ve from_json çıkarabileceği alanlar aynı sorguda başvurulursa analiz aşağıdaki örnekte olduğu gibi başarısız olabilir:
CREATE STREAMING TABLE bronze AS
SELECT
from_json(value, NULL, map('schemaLocationKey', 'nycTaxi')) jsonCol
FROM STREAM READ_FILES('/databricks-datasets/nyctaxi/sample/json/', format => 'text')
WHERE jsonCol.total_amount > 100.0
Bu sorunu, referansı daha sonraki sorguya from_json taşıyarak (yukarıdaki bronz/gümüş örneği gibi) düzeltebilirsiniz. Alternatif olarak, başvurulan schemaHints alanlarını içerecek şekilde from_json belirtmek gibi seçenekler belirtebilirsiniz. Örneğin:
CREATE STREAMING TABLE bronze AS
SELECT
from_json(value, NULL, map('schemaLocationKey', 'nycTaxi', 'schemaHints', 'total_amount DOUBLE')) jsonCol
FROM STREAM READ_FILES('/databricks-datasets/nyctaxi/sample/json/', format => 'text')
WHERE jsonCol.total_amount > 100.0
Örnekler: Şemayı otomatik olarak çıkar ve geliştir
Bu bölümde, işlem hatlarında kullanarak from_json otomatik şema çıkarımı ve evrimi etkinleştirmeye yönelik örnek kod sağlanır.
Bulut nesne depolama alanından akış tablosu oluşturma
Aşağıdaki örnek, read_files söz dizimini kullanarak bulut nesne deposundan bir akış tablosu oluşturur.
SQL
CREATE STREAMING TABLE bronze AS
SELECT
from_json(value, NULL, map('schemaLocationKey', 'nycTaxi')) jsonCol
FROM STREAM READ_FILES('/databricks-datasets/nyctaxi/sample/json/', format => 'text')
Piton
@dp.table(comment="from_json autoloader example")
def bronze():
return (
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "text")
.load("/databricks-datasets/nyctaxi/sample/json/")
.select(from_json(col("value"), None, {"schemaLocationKey": "nycTaxi"}).alias("jsonCol"))
)
Kafka'dan akış tablosu oluşturma
Aşağıdaki örnek, read_kafka sözdizimini kullanarak Kafka'dan bir akış tablosu oluşturur.
SQL
CREATE STREAMING TABLE bronze AS
SELECT
value,
from_json(value, NULL, map('schemaLocationKey', 'keyX')) jsonCol,
FROM READ_KAFKA(
bootstrapSevers => '<server:ip>',
subscribe => 'events',
"startingOffsets", "latest"
)
Piton
@dp.table(comment="from_json kafka example")
def bronze():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "<server:ip>")
.option("subscribe", "<topic>")
.option("startingOffsets", "latest")
.load()
.select(col(“value”), from_json(col(“value”), None, {"schemaLocationKey": "keyX"}).alias("jsonCol"))
)
Örnekler: Sabit şema
Örneğin, sabit bir şema ile kullanılan from_json örnek kodu için bkz from_json işlevi.
FAQs
Bu bölüm, işlevdeki şema çıkarımı ve evrim desteği hakkında sık sorulan soruları yanıtlar from_json .
ile from_jsonarasındaki parse_json fark nedir?
İşlev, parse_json JSON dizesinden bir VARIANT değer döndürür.
VARIANT, yarı yapılandırılmış verileri depolamak için esnek ve verimli bir yol sağlar. Bu, katı türleri tamamen ortadan kaldırarak şema çıkarımı ve evrimi atlatır. Ancak, bir şemayı yazma zamanında zorlamak istiyorsanız (örneğin, görece katı bir şemanız olduğundan), from_json daha iyi bir seçenek olabilir.
Aşağıdaki tabloda ile from_jsonarasındaki parse_json farklar açıklanmaktadır:
| İşlev | Kullanım örnekleri | Availability |
|---|---|---|
from_json |
ile from_json şema evrimi, şemayı korur. Bu, aşağıdaki durumlarda yararlı olur:
|
Yalnızca işlem hatlarında şema çıkarımı ve evrimi ile kullanılabilir |
parse_json |
VARIANT, şemalanması gerekmeyen verileri tutmak için özellikle uygundur. Örneğin:
|
İşlem hatlarının içinde ve dışında kullanılabilir |
from_json şema çıkarımı ve evrim söz dizimini işlem hatları dışında kullanabilir miyim?
Hayır, işlem hatlarının dışında şema çıkarımını ve evrim söz dizimini kullanamazsınız from_json .
tarafından çıkarsanan from_jsonşemaya nasıl erişebilirim?
Hedef akış tablosunun şemasını görüntüleyin.
Bir şema geçirip from_json evrim geçirebilir miyim?
Hayır, bir şema geçirip from_json evrim geçiremezsiniz. Ancak, from_json tarafından çıkarsanan öğelerin bazılarını veya tümünü geçersiz kılmak için şema ipuçları sağlayabilirsiniz.
Tablo tamamen yenilenirse şemaya ne olur?
Tabloyla ilişkili şema konumları temizlenir ve şema sıfırdan yeniden oluşturulur.