Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
В этом руководстве показано, как разрабатывать и развертывать первый конвейер ETL (извлечение, преобразование и загрузку) для оркестрации данных с помощью Apache Spark. Хотя в этом руководстве используются универсальные вычислительные ресурсы Databricks, вы также можете использовать бессерверные вычисления, если они включены для вашей рабочей области.
Конвейеры Lakeflow также можно использовать для создания конвейеров ETL. Конвейеры Lakeflow снижают сложность создания, развертывания и обслуживания рабочих конвейеров ETL. См. руководство по созданию конвейера ETL с помощью конвейеров Lakeflow.
К концу этой статьи вы узнаете, как:
- Запустите вычислительный ресурс Databricks для всех целей.
- Создайте записную книжку Databricks.
- Настройте добавочное прием данных в Delta Lake с помощью автозагрузчика.
- Обработка и взаимодействие с данными.
- Запланируйте записную книжку в качестве задания Databricks.
В этом руководстве используются интерактивные записные книжки для выполнения распространенных задач ETL в Python или Scala.
Вы также можете использовать поставщик Databricks Terraform для создания ресурсов этой статьи. См. статью "Создание кластеров, записных книжек и заданий с помощью Terraform".
Требования
- Вы вошли в рабочую область Azure Databricks.
- У вас есть разрешение на создание вычислительного ресурса.
Примечание.
Если у вас нет прав управления вычислениями, вы все равно можете выполнить большую часть приведенных ниже шагов, пока у вас есть доступ к вычислительному ресурсу.
Шаг 1. Создание вычислительного ресурса
Чтобы выполнить анализ и проектирование данных, создайте вычислительный ресурс для выполнения команд.
Примечание.
Если рабочая область включена для бессерверных вычислений, этот шаг можно пропустить. Блокноты автоматически подключаются к бессерверным вычислительным ресурсам при выполнении кода, или вы можете выбрать Serverless в раскрывающемся списке вычислительных ресурсов. См. Бессерверные вычисления для блокнотов.
- Щелкните "
Вычисления" на боковой панели. - На странице вычислений нажмите кнопку "Создать вычисления".
- Укажите уникальное имя вычислительного ресурса, оставьте оставшиеся значения в состоянии по умолчанию и нажмите кнопку "Создать вычисления".
Дополнительные сведения о вычислениях Databricks см. в разделе "Вычисления".
Шаг 2. Создание записной книжки Databricks
Чтобы создать записную книжку в рабочей области, нажмите кнопку
"Создать" на боковой панели и нажмите кнопку "Записная книжка". Пустая записная книжка открывается в рабочей области.
Дополнительные сведения о создании записных книжек и управлении ими см. в статье "Управление записными книжками Databricks".
Шаг 3. Настройка автозагрузчика для приема данных в Delta Lake
В Databricks рекомендуется использовать Автозагрузчик для добавочного приема данных. Автозагрузчик автоматически определяет и обрабатывает новые файл, когда они поступают в облачное объектное хранилище.
Databricks рекомендует хранить данные с использованием формата Delta Lake. Delta Lake — это уровень хранения с открытым исходным кодом, предоставляющий транзакции ACID и обеспечивающий озерохранилище данных. Delta Lake — это формат по умолчанию для таблиц, созданных в Databricks.
Чтобы настроить автозагрузчик для приема данных в таблицу Delta Lake, скопируйте следующий код в пустую ячейку записной книжки:
Python
# Import functions
from pyspark.sql.functions import col, current_timestamp
# Define variables used in code below
file_path = "/databricks-datasets/structured-streaming/events"
username = spark.sql("SELECT regexp_replace(session_user(), '[^a-zA-Z0-9]', '_')").first()[0]
table_name = f"{username}_etl_quickstart"
checkpoint_path = f"/tmp/{username}/_checkpoint/etl_quickstart"
# Clear out data from previous demo execution
spark.sql(f"DROP TABLE IF EXISTS {table_name}")
dbutils.fs.rm(checkpoint_path, True)
# Configure Auto Loader to ingest JSON data to a Delta table
(spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaLocation", checkpoint_path)
.load(file_path)
.select("*", col("_metadata.file_path").alias("source_file"), current_timestamp().alias("processing_time"))
.writeStream
.option("checkpointLocation", checkpoint_path)
.trigger(availableNow=True)
.toTable(table_name))
язык программирования Scala
// Imports
import org.apache.spark.sql.functions.current_timestamp
import org.apache.spark.sql.streaming.Trigger
import spark.implicits._
// Define variables used in code below
val file_path = "/databricks-datasets/structured-streaming/events"
val username = spark.sql("SELECT regexp_replace(session_user(), '[^a-zA-Z0-9]', '_')").first.get(0)
val table_name = s"${username}_etl_quickstart"
val checkpoint_path = s"/tmp/${username}/_checkpoint"
// Clear out data from previous demo execution
spark.sql(s"DROP TABLE IF EXISTS ${table_name}")
dbutils.fs.rm(checkpoint_path, true)
// Configure Auto Loader to ingest JSON data to a Delta table
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaLocation", checkpoint_path)
.load(file_path)
.select($"*", $"_metadata.file_path".as("source_file"), current_timestamp.as("processing_time"))
.writeStream
.option("checkpointLocation", checkpoint_path)
.trigger(Trigger.AvailableNow)
.toTable(table_name)
Примечание.
Переменные, определенные в этом коде, должны позволить безопасно выполнять его без риска конфликта с существующими ресурсами рабочей области или другими пользователями. Ограниченные разрешения сети или хранилища приведут к ошибкам при выполнении этого кода; обратитесь к администратору рабочей области, чтобы решить проблему этих ограничений.
Дополнительные сведения об автозагрузчике см. в статье Автозагрузчик.
Шаг 4. Обработка данных и взаимодействие с ними
Ноутбуки выполняют ячейки логики поэтапно. Для выполнения логики в ячейке:
Чтобы запустить ячейку, выполненную на предыдущем шаге, выберите ячейку и нажмите клавиши SHIFT+ВВОД.
Чтобы запросить только что созданную таблицу, скопируйте и вставьте следующий код в пустую ячейку, а затем нажмите клавиши SHIFT+ВВОД , чтобы запустить ячейку.
Python
df = spark.read.table(table_name)язык программирования Scala
val df = spark.read.table(table_name)Чтобы просмотреть данные в только что созданной таблице, скопируйте следующий код в пустую ячейку, а затем нажмите клавиши SHIFT+ВВОД, чтобы выполнить код в этой ячейке.
Python
display(df)язык программирования Scala
display(df)
Дополнительные сведения об интерактивных параметрах визуализации данных см. в разделе "Визуализации" в записных книжках Databricks и редакторе SQL.
Шаг 5. Планирование задания
Записные книжки Databricks можно запускать в качестве рабочих сценариев, добавляя их в виде задачи в задание Databricks. В этом шаге вы создадите новое задание, которое можно активировать вручную.
Примечание.
Если вы используете бессерверные вычисления, выберите "Бессерверные " в раскрывающемся списке "Вычисления " вместо вычислительного ресурса на шаге 1.
Чтобы запланировать выполнение записной книжки в качестве задачи, выполните следующие действия.
- Нажмите Расписание справа от строки заголовка.
- Введите уникальное Имя задания.
- Выберите Вручную.
- В раскрывающемся списке вычислений выберите вычислительный ресурс, созданный на шаге 1.
- Нажмите кнопку Создать.
- В открывшемся окне нажмите Запустить сейчас.
- Чтобы просмотреть результаты выполнения задания, щелкните на значок
, находящийся рядом с меткой времени последнего выполнения.
Дополнительные сведения о заданиях см. в разделе Что такое задания?.
Дополнительные интеграции
Дополнительные сведения об интеграции и средствах для проектирования данных с помощью Azure Databricks: