Táblák áthelyezése csővezetékek között

Helyezze át a streamingtáblákat és a materializált nézeteket a folyamatok között, hogy az eredeti helyett a célfolyamat frissítse őket. Ez számos forgatókönyvben hasznos, például:

  • Nagy csővezeték felosztása kisebbekre.
  • Több csővezeték egyesítése egyetlen nagyobban.
  • Néhány táblázat frissítési gyakoriságának módosítása egy folyamatban.
  • Táblák áthelyezése az örökölt közzétételi módot használó folyamatból az alapértelmezett közzétételi módba. Az örökölt közzétételi móddal kapcsolatos részletekért tekintse meg a csővezetékek örökölt közzétételi módját. Ha meg szeretné tudni, hogyan migrálhatja egyszerre egy teljes folyamat közzétételi módját, olvassa el az Alapértelmezett közzétételi mód engedélyezése folyamatban című témakört.
  • Táblák áthelyezése különböző munkaterületek pipeline-jei között.

Requirements

A következő feltételek vonatkoznak egy táblázat áthelyezésére csővezetékek között.

  • A parancs futtatásakor a Databricks Runtime 16.3-at és a Databricks Runtime 17.2-et kell használnia a ALTER ... munkaterületek közötti tábla áthelyezéséhez.

  • A forrás- és célfolyamatoknak metaadattárat használó munkaterületeken kell lenniük. A metaadattár ellenőrzéséhez lásd current_metastore a függvényt.

  • A műveletet futtató felhasználói fióknak vagy szolgáltatásnévnek a forrás- és célfolyamatok futtató felhasználójaként kell lennie.

  • Önnek kell a forrás- és a célfolyamat tulajdonosának lennie.

  • A célfolyamatnak az alapértelmezett közzétételi módot kell használnia. Így több katalógusban és sémában is közzéteheti a táblákat.

    Másik megoldásként mindkét folyamatnak az örökölt közzétételi módot kell használnia, és mindkettőnek ugyanazzal a katalógussal és célértékel kell rendelkeznie a beállításokban. Az örökölt közzétételi módról további információt a LIVE séma (örökölt) című témakörben talál.

    Az örökölt közzétételi módú folyamatok a folyamatbeállítások felhasználói felületének Összegzés mezőjében jelennek meg.

Megjegyzés:

Ez a funkció nem támogatja az alapértelmezett közzétételi módot használó folyamat áthelyezését egy folyamatba az örökölt közzétételi móddal.

Táblázat áthelyezése csővezetékek között

Az alábbi utasítások bemutatják, hogyan helyezhet át streamelési táblázatot vagy materializált nézetet egyik folyamatból a másikba.

  1. Állítsa le a forráscsővezetéket, ha éppen fut. Várjon, amíg teljesen leáll.

  2. Távolítsa el a tábla definícióját a forrásfolyamat kódjából, és tárolja valahol későbbi referenciaként.

    Adjon meg minden olyan támogató lekérdezést vagy kódot, amely a folyamat megfelelő futtatásához szükséges.

  3. Jegyzetfüzetből vagy SQL-szerkesztőből futtassa a következő SQL-parancsot a tábla forrásfolyamatból a célfolyamathoz való újbóli hozzárendeléséhez:

    ALTER [MATERIALIZED VIEW | STREAMING TABLE | TABLE] <table-name>
      SET TBLPROPERTIES("pipelines.pipelineId"="<destination-pipeline-id>");
    

    Vegye figyelembe, hogy az SQL-parancsot a forrásfolyamat munkaterületéről kell futtatni.

    A parancs a Unity Catalog által felügyelt materializált nézetekhez ALTER MATERIALIZED VIEW-t, a streamtáblákhoz pedig ALTER STREAMING TABLE-t használ. Ha ugyanezt a műveletet egy Hive metaadattártáblán szeretné végrehajtani, használja a következőt ALTER TABLE: .

    Például, ha egy sales nevű streamelési táblát szeretne áthelyezni egy abcd1234-ef56-ab78-cd90-1234efab5678 azonosítóval rendelkező folyamatba, akkor a következő parancsot kell futtatnia:

    ALTER STREAMING TABLE sales
      SET TBLPROPERTIES("pipelines.pipelineId"="abcd1234-ef56-ab78-cd90-1234efab5678");
    

    Megjegyzés:

    A pipelineId folyamatazonosítónak érvényesnek kell lennie. Az null érték nem engedélyezett.

  4. Adja hozzá a tábla definícióját a célfolyamat kódjába.

    Megjegyzés:

    Ha a katalógus vagy a célséma eltér a forrás és a cél között, előfordulhat, hogy a lekérdezés másolása nem működik. A részben minősített táblák a definícióban eltérően oldhatók fel. Előfordulhat, hogy frissítenie kell a definíciót a táblanevek teljes minősítéséhez.

    Megjegyzés:

    Távolítsa el vagy kommentezze ki az egyszeri hozzáfűzési folyamatokat a cél pipeline kódjából (Pythonban, append_flow(once=True), SQL-ben, INSERTINTO ONCE lekérdezések). További részletekért lásd: Korlátozások.

Az áthelyezés befejeződött. Most már futtathatja a forrás- és célfolyamatokat is. A célfolyamat frissíti a táblát.

Hibaelhárítás

Az alábbi táblázat azokat a hibákat ismerteti, amelyek akkor fordulhatnak elő, ha egy táblát a folyamatcsatornák között áthelyez.

Error Description
DESTINATION_PIPELINE_NOT_IN_DIRECT_PUBLISHING_MODE A forrásfolyamat az alapértelmezett közzétételi módban van, és a cél a LIVE séma (örökölt) üzemmódot használja. Ez nem támogatott. Ha a forrás az alapértelmezett közzétételi módot használja, akkor a célhelynek is meg kell jelennie.
PIPELINE_TYPE_NOT_WORKSPACE_PIPELINE_TYPE Csak a táblák adatfolyamok közötti áthelyezése támogatott. Az önálló streamelési táblák és a materializált nézetek áthelyezése nem támogatott.
DESTINATION_PIPELINE_NOT_FOUND A pipelines.pipelineId csővezetéknek érvényesnek kell lennie. A pipelineId nem lehet null értékű.
A tábla az áthelyezés után nem frissül a célhelyen. Ebben az esetben a gyors enyhítés érdekében helyezze vissza a táblát a forrásfolyamatba ugyanezeket az utasításokat követve.
PIPELINE_PERMISSION_DENIED_NOT_OWNER A forrás- és célfolyamatoknak egyaránt az áthelyezési műveletet végrehajtó felhasználó tulajdonában kell lenniük.
TABLE_ALREADY_EXISTS A hibaüzenetben felsorolt tábla már létezik. Ez akkor fordulhat elő, ha az adatcsatornához már létezik támogató tábla. Ebben az esetben DROP a hiba által érintett táblázat.

Példa több táblára egy folyamatábrában

A csővezetékek több táblát is tartalmazhatnak. A vizorendszerek között egyszerre egy táblát is áthelyezhet. Ebben a forgatókönyvben három tábla (table_a, table_b, table_c) olvasható egymástól egymás után a forrásfolyamatban. Át szeretnénk helyezni egy táblát, table_b, egy másik folyamatba.

Kezdeti forrásfolyamat kódja:

from pyspark import pipelines as dp
from pyspark.sql.functions import col, sum

@dp.table
def table_a():
    return spark.read.table("source_table")

# Table to be moved to new pipeline:
@dp.table
def table_b():
    return (
        spark.read.table("table_a")
        .select(col("column1"), col("column2"))
    )

@dp.table
def table_c():
    return (
        spark.read.table("table_b")
        .groupBy(col("column1"))
        .agg(sum("column2").alias("sum_column2"))
    )

Átmozgatjuk table_b egy másik folyamatba, a forrásból a tábladefiníció másolásával és eltávolításával, valamint table_b pipelineId-jának frissítésével.

Először szüneteltetje az ütemezéseket, és várja meg, amíg a frissítések befejeződnek a forrás- és a célfolyamatokon. Ezután módosítsa a forrásfolyamatot az áthelyezett tábla kódjának eltávolításához. A frissített forrásfolyamat példakódja a következő lesz:

from pyspark import pipelines as dp
from pyspark.sql.functions import col, sum

@dp.table
def table_a():
    return spark.read.table("source_table")

# Removed, to be in new pipeline:
# @dp.table
# def table_b():
#     return (
#         spark.read.table("table_a")
#         .select(col("column1"), col("column2"))
#     )

@dp.table
def table_c():
    return (
        spark.read.table("table_b")
        .groupBy(col("column1"))
        .agg(sum("column2").alias("sum_column2"))
    )

Nyissa meg az SQL-szerkesztőt a ALTER pipelineId parancs futtatásához.

ALTER MATERIALIZED VIEW table_b
  SET TBLPROPERTIES("pipelines.pipelineId"="<new-pipeline-id>");

Ezután lépjen a célcsővezetékhez, és adja hozzá a table_b definíciót. Ha az alapértelmezett katalógus és séma megegyezik a folyamat beállításai között, nincs szükség kódmódosításra.

A célfolyamat kódja:

from pyspark import pipelines as dp
from pyspark.sql.functions import col

@dp.table(name="table_b")
def table_b():
    return (
        spark.read.table("table_a")
        .select(col("column1"), col("column2"))
    )

Ha az alapértelmezett katalógus és séma eltér a folyamat beállításaiban, a teljesen kvalifikált nevet a folyamat katalógusának és sémájának használatával kell hozzáadnia.

A célfolyamat kódja például a következő lehet:

from pyspark import pipelines as dp
from pyspark.sql.functions import col

@dp.table(name="source_catalog.source_schema.table_b")
def table_b():
    return (
        spark.read.table("source_catalog.source_schema.table_a")
        .select(col("column1"), col("column2"))
    )

Futtassa (vagy engedélyezze újra az ütemezéseket) a forrás- és a célfolyamatokhoz is.

A csővezetékek most már különállóak. A table_c a table_b-ből olvas (most a célfolyamatban), és a table_b a table_a-ból olvas (a forrásfolyamatban). Ha indított végrehajtást hajt végre a forrásfolyamaton, table_b nem frissül, mert azt már a forrásfolyamat nem kezeli. A forrásfolyamat a folyamaton kívüli táblaként kezeli table_b . Ez hasonló ahhoz, hogy a folyamat által nem felügyelt Unity-katalógus egyik Delta-táblájából származó materializált nézetet definiáljon.

korlátozások

Az alábbiakban bemutatjuk a táblák munkafolyamatok közötti áthelyezésének korlátozásait.

  • Az önálló materializált nézetek és streamtáblák nem támogatottak.
  • A hozzáfűző egyszeri folyamatok – Python append_flow(once=True) és SQL INSERT INTO ONCE folyamatok – nem támogatottak. A futtatási állapotuk nem marad meg, és előfordulhat, hogy újra futnak a célfolyamatban. Távolítsa el vagy kommentálja ki az egyszer történő csatolás folyamatait a célfolyamatból, hogy elkerülje ezek újbóli futtatását.
  • A privát táblák vagy nézetek nem támogatott.
  • A forrás- és célcsővezetékeknek csővezetékeknek kell lenniük. A null értékű csővezetékek nem támogatottak.
  • A forrás- és célfolyamatoknak ugyanabban a munkaterületen vagy különböző munkaterületeken kell lenniük, amelyek ugyanazt a metaadattárat használják.
  • A forrás- és célfolyamatoknak az áthelyezési műveletet futtató felhasználó tulajdonában kell lenniük.
  • Ha a forrásfolyamat az alapértelmezett közzétételi módot használja, a célfolyamatnak az alapértelmezett közzétételi módot is használnia kell. Nem helyezhet át táblázatot egy folyamatból az alapértelmezett közzétételi móddal olyan folyamatba, amely a LIVE sémát (örökölt) használja. Lásd: LIVE séma (régi változat).
  • Ha a forrás- és célfolyamatok a LIVE sémát (örökölt) használják, akkor a beállításokban ugyanazokkal catalog és target értékekkel kell rendelkezniük.