İşlem hatlarında kullanarak from_json şemayı çıkarsama ve geliştirme

Ö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. buradadaha fazla bilgi bulabilirsiniz.

Ş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:
  • Veri şemanızı zorunlu kılmak istiyorsunuz (örneğin, kalıcı hale getirmeden önce her şema değişikliğini gözden geçirme).
  • Depolamayı iyileştirmek ve düşük sorgu gecikme süresi ve maliyet gerektirmek istiyorsunuz.
  • Eşleşmeyen türlerdeki verilerde başarısız olmak istiyorsunuz.
  • Bozuk JSON kayıtlarından kısmi sonuçları ayıklamak ve bozuk kaydı sütun _corrupt_record içinde depolamak istiyorsunuz. Buna karşılık, VARIANT alımı geçersiz JSON için bir hata döndürür.
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:
  • Esnek olduğundan verileri yarı yapılandırılmış tutmak istiyorsunuz.
  • Şema, sık sık akış hataları ve yeniden başlatmalar olmadan şemaya dönüştürülemeyecek kadar hızlı değişir.
  • Eşleşmeyen türlerdeki verilerde başarısız olmak istemezsiniz. (VARIANT içeri alma işlemi, tür uyuşmazlıkları olsa bile geçerli JSON kayıtlarında her zaman başarılı olur.)
  • Kullanıcılarınız, şemaya uymayan alanlar içeren kurtarılan veri sütunuyla ilgilenmek istemiyor.
İş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.