Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
На этой странице перечислены доступные параметры входных и выходных данных для API Spark, которые считывают и записывают данные.
Параметры DataFrameReader
Используйте эти параметры с DataFrameReader.option(), DataFrameReader.options(), read_files, COPY INTO и Auto Loader для управления тем, как Azure Databricks считывает файлы данных.
Example
В следующем примере задано multiLine значение True для чтения JSON-файлов:
Python
df = spark.read.format("json").option("multiLine", True).load("/path/to/data")
Scala
val df = spark.read.format("json").option("multiLine", "true").load("/path/to/data")
SQL
SELECT * FROM read_files("/path/to/data", format => "json", multiLine => true)
Общие
Следующие параметры применяются ко всем форматам файлов.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
ignoreCorruptFiles |
false |
true, false |
Определяет, следует ли игнорировать поврежденные файлы. Если задано значение true, задания Spark будут продолжать выполняться при обнаружении поврежденных файлов, а прочитанное содержимое будет возвращено. Для COPY INTOэтого можно наблюдать пропущенные поврежденные файлы, как numSkippedCorruptFiles в столбце operationMetrics журнала Delta Lake. Доступно в Databricks Runtime 11.3 LTS и более поздних версиях. |
ignoreMissingFiles |
false для автозагрузчика trueCOPY INTO (устаревшая версия) |
true, false |
Определяет, следует ли игнорировать отсутствующие файлы. Если задано значение true, задания Spark продолжают выполняться при обнаружении отсутствующих файлов, а содержимое по-прежнему возвращается. Доступно в Databricks Runtime 11.3 LTS и более поздних версиях. |
modifiedAfter |
None | Строка метки времени | Необязательная метка времени в качестве фильтра для приема только файлов с меткой времени изменения после указанной метки времени. |
modifiedBefore |
None | Строка метки времени | Необязательная метка времени в качестве фильтра для приема только файлов с меткой времени изменения до указанной метки времени. |
pathGlobFilter или fileNamePattern |
None | Строка шаблона глоба | Потенциальный шаблон глобов для выбора файлов. Эквивалентно PATTERNCOPY INTO (устаревшей версии).
fileNamePattern можно использовать в read_files. |
recursiveFileLookup |
false |
true, false |
Если trueэтот параметр выполняет поиск по вложенным каталогам, даже если их имена не соответствуют схеме именования секций, например date=2019-07-01. |
Avro
При чтении файлов Avro применяются следующие параметры.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
avroSchema |
None | Строка схемы Avro | Необязательная схема, указанная пользователем в формате Avro. При чтении Avro этот параметр может быть установлен на развивающуюся схему, совместимую, но отличающуюся от фактической схемы Avro. Схема десериализации согласуется с развивающейся схемой. Например, если задать развиваемую схему, содержащую один дополнительный столбец со значением по умолчанию, результат чтения также содержит новый столбец. |
avroSchemaEvolutionMode |
none |
none, restart |
Как обрабатывать эволюцию схемы при использовании реестра схем.
none игнорирует изменения схемы и продолжает задание.
restart вызывает UnknownFieldException при обнаружении изменений схемы и требуется перезапуск задания. |
datetimeRebaseMode |
LEGACY |
EXCEPTION, , LEGACYCORRECTED |
Управляет изменением базы значений DATE и TIMESTAMP между юлианским и пролептическим григорианским календарями. |
enableStableIdentifiersForUnionType |
false |
true, false |
Следует ли использовать стабильные имена полей для типов Avro Union. При включении имена полей типа объединения являются производными от имен типов в нижнем регистре (например, member_int). member_string Вызывает исключение, если два типа идентичны после нижней части. |
mergeSchema |
false |
true, false |
Следует ли определять схему на основе нескольких файлов и объединять схемы каждого файла.
mergeSchema для Avro не ослабляет ограничения типов данных. |
mode |
FAILFAST |
FAILFAST, , PERMISSIVEDROPMALFORMED |
Режим синтаксического анализа для обработки поврежденных записей.
FAILFAST выдает исключение.
PERMISSIVE Задает для неправильно сформированных полей значение NULL.
DROPMALFORMED автоматически сбрасывает плохие записи. |
readerCaseSensitive |
true |
true, false |
Указывает поведение чувствительности к регистру при включении rescuedDataColumn. Если значение true, спасите столбцы данных, имена которых отличаются по регистру от схемы. Если значение false, считывает данные без учета регистра. |
recursiveFieldMaxDepth |
None |
0 до 15 |
Максимальная глубина рекурсии для рекурсивных полей Avro. Установите для 1 усечения всех рекурсивных полей, 2 чтобы разрешить один уровень рекурсии и т. д 15. Если не задано или 0не разрешено рекурсивные поля. |
rescuedDataColumn |
None | Строка имени столбца | Следует ли собирать все данные, которые не могут быть проанализированы из-за несоответствия типов данных и несоответствия схемы (включая регистр столбцов) отдельному столбцу. Этот столбец включен по умолчанию при использовании Автозагрузчика.COPY INTO (устаревшая версия) не поддерживает спасательный столбец данных, так как невозможно вручную задать схему с помощью COPY INTO. Databricks рекомендует использовать Auto Loader для большинства сценариев загрузки.Дополнительные сведения см. в статье "Что такое столбец спасенных данных?". |
stableIdentifierPrefixForUnionType |
member_ |
Любая строка. | Префикс, используемый для стабильных имен полей типа объединения, когда enableStableIdentifiersForUnionType=true. |
CSV
При чтении CSV-файлов применяются следующие параметры.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
badRecordsPath |
None | Строка пути | Путь для хранения файлов, содержащих сведения об ошибочных записях CSV. |
charToEscapeQuoteEscaping |
\0 |
Один символ | Символ, используемый для экранирования того символа, который используется для экранирования кавычек. Например, для следующей записи: [ " a\\", b ]:
|
columnNameOfCorruptRecord |
_corrupt_record |
Строка имени столбца | Поддержка для автозагрузчика. Не поддерживается для COPY INTO (устаревшая версия).Столбец для хранения записей, которые имеют неправильный формат и не могут быть проанализированы. Если в качестве mode для синтаксического анализа задано значение DROPMALFORMED, этот столбец будет пустым. |
comment |
\0 |
Один символ | Определяет символ, который обозначает строку комментария, если он найден в начале строки текста. Используйте '\0', чтобы отключить пропуск комментариев. |
dateFormat |
yyyy-MM-dd |
Строка формата даты | Формат синтаксического анализа строк даты. |
emptyValue |
Пустая строка | Любая строка. | Строковое представление пустого значения. |
enableDateTimeParsingFallback |
false |
true, false |
Следует ли вернуться к устаревшей дате и поведению синтаксического анализа меток времени, если значение не может быть проанализировано с указанным форматом. При falseсинтаксическом анализе возникает ошибка или в зависимости от modeзначения NULL. |
encoding или charset |
UTF-8 |
java.nio.charset.Charset Имя |
Имя кодировки CSV-файлов. Смотрите java.nio.charset.Charset для списка вариантов.
UTF-16 и UTF-32 использовать нельзя, если multiline имеет значение true. |
enforceSchema |
true |
true, false |
Следует ли принудительно применять указанную или выведенную схему к CSV-файлам. Если параметр включен, заголовки CSV-файлов игнорируются. Этот параметр игнорируется по умолчанию при использовании Автозагрузчика для восстановления данных и разрешения на развитие схемы. |
escape |
\ |
Один символ | Escape-символ, используемый при анализе данных. |
extension |
csv |
Строка расширения файла | Ожидаемое расширение имени файла для операций чтения. Файлы без этого расширения отфильтровываются. |
failOnUnknownFields |
false |
true, false |
Происходит ли сбой, если запись CSV содержит столбцы, не присутствующих в схеме. Если falseнераспознанные столбцы автоматически удаляются или спасаются в зависимости от rescuedDataColumnних. |
failOnWidenedFields |
false |
true, false |
Не удается ли выполнить анализ значения поля как объявленного типа схемы без расширения. Если falseзначения с расширением типа автоматически спасаются в зависимости от rescuedDataColumn. Параметр failOnUnknownFields=true может маскировки эффектов этого параметра. |
header |
false |
true, false |
Содержат ли CSV-файлы заголовок. При выводе схемы Автозагрузчик предполагает, что файлы имеют заголовки. |
ignoreLeadingWhiteSpace |
false |
true, false |
Следует ли игнорировать начальные пробелы для каждого анализируемого значения. |
ignoreTrailingWhiteSpace |
false |
true, false |
Следует ли игнорировать завершающие пробелы для каждого анализируемого значения. |
inferSchema |
false |
true, false |
Указывает, следует ли вычислять типы данных проанализированных записей CSV, или предполагается, что все столбцы имеют тип StringType. Требует дополнительного прохода по данным, если задано значение true. Для автозагрузчика используйте cloudFiles.inferColumnTypes вместо этого. |
inputBufferSize |
1048576 (1 МБ) |
Положительные целые числа | Размер буфера в байтах для средства синтаксического анализа CSV. Полезно для настройки использования памяти при анализе больших CSV-файлов. |
lineSep |
Нет, который охватывает \r, \r\nи \n |
Строка | Строка между двумя последовательными записями CSV. |
locale |
US |
Идентификатор java.util.Locale |
Язык Java, который влияет на дату, метку времени и десятичный анализ по умолчанию в CSV. |
maxCharsPerColumn |
-1 |
Положительные целые числа или -1 неограниченные |
Максимальное число символов, ожидаемое в значении для синтаксического анализа. Можно использовать, чтобы избежать ошибок памяти. По умолчанию имеет значение -1, что означает отсутствие ограничений. |
maxColumns |
20480 |
Положительные целые числа | Фиксированное ограничение количества столбцов в записи. |
mergeSchema |
false |
true, false |
Следует ли определять схему на основе нескольких файлов и объединять схемы каждого файла. По умолчанию включено для автозагрузчика при выведении схемы. |
mode |
PERMISSIVE |
PERMISSIVE, , DROPMALFORMEDFAILFAST |
Режим парсера для обработки некорректных записей. |
multiLine |
false |
true, false |
Являются ли записи CSV многострочными. |
nanValue |
NaN |
Любая строка. | Строковое представление значения, не являющегося числовым, при синтаксическом анализе столбцов FloatType и DoubleType. |
negativeInf |
-Inf |
Любая строка. | Строковое представление отрицательной бесконечности при синтаксическом анализе столбцов FloatType или DoubleType. |
nullValue |
Пустая строка | Любая строка. | Строковое представление значения NULL. |
parserCaseSensitive (не рекомендуется) |
false |
true, false |
Следует ли при чтении файлов выравнивать столбцы, указанные в заголовке, с учетом регистра в соответствии со схемой? Это значение равно true по умолчанию для Автозагрузчика. Столбцы, отличающиеся только регистром текста, будут сохранены в rescuedDataColumn, если эта функция включена. Вместо этого параметра рекомендуется использовать readerCaseSensitive. |
positiveInf |
Inf |
Любая строка. | Строковое представление положительной бесконечности при синтаксическом анализе столбцов FloatType или DoubleType. |
preferDate |
true |
true, false |
Пытается интерпретировать строки как дату, а не как отметку времени, когда это возможно. Кроме того, необходимо использовать вывод схемы, включив inferSchema или используя cloudFiles.inferColumnTypes функцию автозагрузчика. |
quote |
" |
Один символ | Символ, используемый для изоляции значений, в которых разделитель полей является частью значения. |
readerCaseSensitive |
true |
true, false |
Указывает поведение чувствительности к регистру при включении rescuedDataColumn. Если значение true, спасите столбцы данных, имена которых отличаются по регистру от схемы. Если значение false, считывает данные без учета регистра. |
rescuedDataColumn |
None | Строка имени столбца | Следует ли собирать все данные, которые не могут быть проанализированы из-за несоответствия типов данных и несоответствия схемы (включая регистр столбцов) отдельному столбцу. Этот столбец включен по умолчанию при использовании Автозагрузчика. Дополнительные сведения см. в статье "Что такое столбец спасенных данных?".COPY INTO (устаревшая версия) не поддерживает спасательный столбец данных, так как невозможно вручную задать схему с помощью COPY INTO. Databricks рекомендует использовать Auto Loader для большинства сценариев загрузки. |
sep или delimiter |
, |
Строка | Разделительная строка между столбцами. |
singleVariantColumn |
None | Строка имени столбца | Если задано имя столбца, считывает всю запись CSV в один VariantType столбец с таким именем, а не анализирует каждое поле в собственный столбец. Требует использования header=true. |
skipRows |
0 |
Положительные целые числа или 0 |
Число строк с начала CSV-файла, которое следует игнорировать, включая закомментированные и пустые строки. Если значение header равно true, заголовок будет первой строкой, которая не пропущена или не закомментирована. |
timeFormat |
HH:mm:ss |
Строка формата времени | Формат для синтаксического анализа значений столбцов TimeType . |
timestampFormat |
yyyy-MM-dd'T'HH:mm:ss[.SSS][XXX] |
Строка формата метки времени | Формат для синтаксического анализа строк меток времени. |
timestampNTZFormat |
yyyy-MM-dd'T'HH:mm:ss[.SSS] |
Строка формата метки времени | Формат синтаксического анализа метки времени без строк часового пояса (TimestampNTZType). |
timeZone |
None |
java.time.ZoneId Строка |
Объект java.time.ZoneId, который следует использовать для анализа меток времени и дат. |
unescapedQuoteHandling |
STOP_AT_DELIMITER |
STOP_AT_CLOSING_QUOTE, , BACK_TO_DELIMITERSTOP_AT_DELIMITER, SKIP_VALUERAISE_ERROR |
Стратегия обработки неэкранированных кавычек. Поведение каждого разрешенного параметра выглядит следующим образом:
|
Excel
Следующие параметры применяются при чтении Excel файлов.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
dataAddress |
None | Строка имени ячейки или листа | Диапазон ячеек для чтения в синтаксисе Excel. Если опущено, считывает все допустимые ячейки из первого листа. Используется SheetName!C5:H10 для чтения диапазона от именованного листа, C5:H10 чтения диапазона от первого листа или SheetName чтения всех данных из определенного листа. |
headerRows |
0 |
0, 1 |
Число начальных строк, используемых в качестве заголовков имени столбца. При dataAddress указании это применяется в диапазоне ячеек. Когда 0имена столбцов создаются автоматически как _c1, _c2и _c3т. д. |
ignoreMissingSheet |
false |
true, false |
Следует ли автоматически пропускать файлы, не содержащие лист, указанный в параметре dataAddress. Если falseфайл отсутствует запрошенный лист, возникает ошибка. Применяется только при указании dataAddressимени листа. |
includePhoneticRuns |
false |
true, false |
Следует ли включать фонетические заметки (например, pinyin или furigana), сцепленные со значениями строк ячейки при чтении файлов XLSX. |
operation |
readSheet |
readSheet, listSheets |
Операция, выполняемая в книге Excel.
readSheet считывает данные из листа.
listSheets возвращает структуру с полями sheetIndex: long и sheetName: String для каждого листа. |
timestampNTZFormat |
yyyy-MM-dd'T'HH:mm:ss[.SSS] |
Строка формата метки времени | Строка настраиваемого формата для значений метки времени без часового пояса, хранящихся в виде строк в Excel. Пользовательские форматы даты следуют форматам из шаблонов даты и времени. |
dateFormat |
yyyy-MM-dd |
Строка формата даты | Строка настраиваемого формата для строковых значений, прочитанных как Date. Пользовательские форматы даты следуют форматам из шаблонов даты и времени. |
JSON
При чтении JSON-файлов применяются следующие параметры.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
allowBackslashEscapingAnyCharacter |
false |
true, false |
Разрешить ли использовать обратные косые черты для экранирования любого последующего символа. Если параметр не активирован, экранировать можно только те символы, которые эксплицитно указаны в спецификации JSON. |
allowComments |
false |
true, false |
Разрешить ли использование комментариев в стиле Java, C и C++ (видов '/', '*' и '//') в проанализированном содержимом. |
allowNonNumericNumbers |
true |
true, false |
Разрешить ли использование набора токенов, представляющих нечисловые значения (NaN), как законные числовые значения с плавающей запятой. |
allowNumericLeadingZeros |
false |
true, false |
Разрешить ли целочисленное число начинаться с дополнительных (игнорируемых) нули (например, 000001). |
allowSingleQuotes |
true |
true, false |
Разрешить ли использование одинарных кавычек (апостроф, символ '\') для цитирования строк (имен и строковых значений). |
allowUnquotedControlChars |
false |
true, false |
Следует ли разрешить строкам JSON содержать управляющие символы без экранирования (символы ASCII со значением менее 32, включая символы табуляции и перевода строки) или нет. |
allowUnquotedFieldNames |
false |
true, false |
Следует ли разрешать использование неквотируемых имен полей, разрешенных JavaScript, но не спецификацией JSON. |
alternateVariantEncoding |
None | Z85 |
Кодировка, используемая для значений Variant в исходном формате JSON. Задайте для Z85 декодирования значений Variant, которые были закодированы в Кодировке Base85 вместо встроенного JSON. |
badRecordsPath |
None | Строка пути | Путь для хранения файлов для записи сведений о неправильных записях JSON.badRecordsPath Использование параметра в источнике данных на основе файлов имеет следующие ограничения:
|
columnNameOfCorruptRecord |
_corrupt_record |
Строка имени столбца | Столбец для хранения записей, которые имеют неправильный формат и не могут быть проанализированы. Если в качестве mode для синтаксического анализа задано значение DROPMALFORMED, этот столбец будет пустым. |
dateFormat |
yyyy-MM-dd |
Строка формата даты | Формат синтаксического анализа строк даты. |
dropFieldIfAllNull |
false |
true, false |
Игнорировать ли столбцы, в которых все значения равны NULL, пустые массивы и структуры во время вывода схемы. |
encoding или charset |
UTF-8 |
java.nio.charset.Charset Имя |
Имя кодировки файлов JSON. Список вариантов см. в java.nio.charset.Charset. Нельзя использовать UTF-16 и UTF-32, если multiline имеет значение true. |
inferTimestamp |
false |
true, false |
Следует ли попытаться распознать строки меток времени как TimestampType. Если задано значение true, вывод схемы может занять заметно больше времени. Необходимо включить cloudFiles.inferColumnTypes, чтобы использовать автозагрузчик. |
lineSep |
Нет, который охватывает \r, \r\nи \n |
Строка | Строка между двумя последовательными записями JSON. |
locale |
US |
Идентификатор java.util.Locale |
Java идентификатор языкового стандарта, влияющий на дату, метку времени и десятичный анализ в формате JSON. |
maxNestingDepth |
500 |
Положительные целые числа | Максимальная разрешенная глубина вложения для объектов и массивов JSON. Увеличьте это значение для глубоко вложенных документов. |
maxNumLen |
1000 |
Положительные целые числа | Максимальная длина маркеров числа в входных данных JSON. Увеличьте это значение для JSON с большими числовыми литералами. |
maxStringLen |
неограниченно | Положительные целые числа | Максимальная длина строковых значений во входных данных JSON. Установите для ограничения использования памяти при анализе JSON с большими строками. |
mode |
PERMISSIVE |
PERMISSIVE, , DROPMALFORMEDFAILFAST |
Режим парсера для обработки некорректных записей. |
multiLine |
false |
true, false |
Занимает ли запись JSON несколько строк. |
prefersDecimal |
false |
true, false |
Пытается интерпретировать строки как DecimalType вместо типа float или double, когда это возможно. Кроме того, необходимо использовать вывод схемы, включив inferSchema или используя cloudFiles.inferColumnTypes функцию автозагрузчика. |
primitivesAsString |
false |
true, false |
Следует ли интерпретировать примитивные типы, такие как числа и логические значения, как StringType. |
readerCaseSensitive |
true |
true, false |
Указывает поведение чувствительности к регистру при включении rescuedDataColumn. Если значение true, спасите столбцы данных, имена которых отличаются по регистру от схемы. Если значение false, считывает данные без учета регистра. Доступно в Databricks Runtime 13.3 и более поздних версиях. |
rescuedDataColumn |
None | Строка имени столбца | Следует ли собирать все данные, которые нельзя проанализировать из-за несоответствия типа данных или несоответствия схемы (включая регистр столбцов) отдельному столбцу. Этот столбец включен по умолчанию при использовании Автозагрузчика. Дополнительные сведения см. в разделе "Что такое столбец спасенных данных?".COPY INTO (устаревшая версия) не поддерживает спасательный столбец данных, так как невозможно вручную задать схему с помощью COPY INTO. Databricks рекомендует использовать Auto Loader для большинства сценариев загрузки. |
singleVariantColumn |
None | Строка имени столбца | Следует ли принять весь документ JSON, проанализированный в один столбец Variant с указанной строкой в качестве имени столбца. Если не задано, поля JSON обрабатываются в собственные столбцы. |
timestampFormat |
yyyy-MM-dd'T'HH:mm:ss[.SSS][XXX] |
Строка формата метки времени | Формат для синтаксического анализа строк меток времени. |
timestampNTZFormat |
yyyy-MM-dd'T'HH:mm:ss[.SSS] |
Строка формата метки времени | Формат синтаксического анализа метки времени без строк часового пояса (TimestampNTZType). |
timeZone |
None |
java.time.ZoneId Строка |
Объект java.time.ZoneId, который следует использовать для анализа меток времени и дат. |
upgradeExceptionAsBadRecord |
false |
true, false |
Следует ли обрабатывать исключения обновления типов (например, если значение не может быть расширено до объявленного типа столбца) как плохие записи, а не исключение. |
Кафка
Полный список параметров чтения Kafka см. в параметрах DataStreamReader Kafka. Следующие параметры применяются только к пакетным считываниям с помощью spark.read.format("kafka").
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
endingOffsets |
latest |
latestили строка смещения JSON |
Где остановить чтение. В строке -1 JSON используется последнее смещение.
-2, который является самым ранним смещением, не допускается в качестве конечного смещения. Это пример строки смещения JSON: {"topicA":{"0":50,"1":-1}} |
endingOffsetsByTimestamp |
None | Строка метки времени JSON | Конечные смещения на секцию, указанные как метки времени в миллисекундах. Например: {"topicA":{"0":2000,"1":3000}}. |
endingTimestamp |
None | Положительные целые числа или 0 |
Глобальная метка времени окончания в миллисекундах, примененная ко всем секциям. |
ОРК
При чтении файлов ORC применяются следующие параметры.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
mergeSchema |
false |
true, false |
Следует ли определять схему на основе нескольких файлов и объединять схемы каждого файла. |
Паркетным
При чтении файлов Parquet применяются следующие параметры.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
datetimeRebaseMode |
LEGACY |
EXCEPTION, , LEGACYCORRECTED |
Управляет изменением базы значений DATE и TIMESTAMP между юлианским и пролептическим григорианским календарями. |
int96RebaseMode |
LEGACY |
EXCEPTION, , LEGACYCORRECTED |
Управляет пересчетом значений временной метки INT96 между юлианским и пролептическим григорианским календарями. |
mergeSchema |
false |
true, false |
Следует ли определять схему на основе нескольких файлов и объединять схемы каждого файла. |
readerCaseSensitive |
true |
true, false |
Указывает поведение чувствительности к регистру при включении rescuedDataColumn. Если значение true, спасите столбцы данных, имена которых отличаются по регистру от схемы. Если значение false, считывает данные без учета регистра. |
rescuedDataColumn |
None | Строка имени столбца | Следует ли собирать все данные, которые не могут быть проанализированы из-за несоответствия типов данных и несоответствия схемы (включая регистр столбцов) отдельному столбцу. Этот столбец включен по умолчанию при использовании Автозагрузчика. Дополнительные сведения см. в статье "Что такое столбец спасенных данных?".COPY INTO (устаревшая версия) не поддерживает спасательный столбец данных, так как невозможно вручную задать схему с помощью COPY INTO. Databricks рекомендует использовать Auto Loader для большинства сценариев загрузки. |
Хранилище состояний
Используйте эти параметры с функцией с spark.read.format("statestore") табличным значением для read_statestore чтения данных о состоянии структурированной потоковой передачи. См. Сведения о состоянии структурированной потоковой передачи.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
batchId |
Последний идентификатор пакетной службы | Положительные целые числа или 0 |
Целевой пакет для чтения. Используется для запроса предыдущего состояния запроса. Пакет должен быть зафиксирован, но еще не очищен. |
operatorId |
0 |
Положительные целые числа или 0 |
Целевой оператор для чтения. Используется, если запрос содержит несколько операторов с отслеживанием состояния. |
storeName |
DEFAULT |
Любая строка. | Имя целевого хранилища состояний для чтения. Используется, когда оператор с отслеживанием состояния содержит несколько экземпляров хранилища состояний. Необходимо указать либо storeName или joinSide для соединения потокового потока, но не для обоих. |
joinSide |
None |
left, right |
Целевая сторона, считываемая из соединения потокового потока. Необходимо указать либо storeName или joinSide для соединения потокового потока, но не для обоих. |
snapshotStartBatchId |
None | Положительные целые числа или 0 |
Идентификатор пакета моментального снимка, используемый в качестве отправной точки при чтении состояния. Средство чтения перестраивает состояние путем повторения изменений из этого моментального снимка до тех пор batchId. Полезно при повреждении моментального снимка. Необходимо указать вместе с snapshotPartitionId. Не удается использовать с readChangeFeed. Поддерживает хранилище состояний, поддерживаемое HDFS, и хранилище состояний RocksDB с включенной контрольной точкой журнала изменений. Доступно в Databricks Runtime 15.4 LTS и более поздних версиях. |
snapshotPartitionId |
None | Положительные целые числа или 0 |
Если задано, запрос считывает только эту секцию. Необходимо указать вместе с snapshotStartBatchId. Не удается использовать с readChangeFeed. Доступно в Databricks Runtime 15.4 LTS и более поздних версиях. |
readChangeFeed |
false |
true, false |
Когда trueвозвращает изменения состояния в заданном диапазоне пакетов между changeStartBatchId и changeEndBatchId. Требует использования changeStartBatchId. Не удается использовать с joinSide, или batchIdsnapshotStartBatchIdsnapshotPartitionId. Доступно в Databricks Runtime 16.4 LTS и более поздних версиях.Дополнительные сведения см. в статье Об изменениях состояния структурированной потоковой передачи. |
changeStartBatchId |
None | Положительные целые числа или 0 |
Начальный идентификатор пакета для диапазона канала изменений. Обязательный, если readChangeFeed имеет значение true. Применяется только в том случае, если readChangeFeed задано значение true. Доступно в Databricks Runtime 16.4 LTS и более поздних версиях. |
changeEndBatchId |
Последний идентификатор пакетной службы | Положительные целые числа или 0 |
Конечный идентификатор пакетной службы для диапазона канала изменений. Должно быть больше или равно changeStartBatchId. Применяется только в том случае, если readChangeFeed задано значение true. Доступно в Databricks Runtime 16.4 LTS и более поздних версиях. |
stateVarName |
None | Любая строка. | Имя переменной состояния для чтения. Имя переменной состояния — это уникальное имя каждой переменной в init функции оператора, используемой StatefulProcessortransformWithState оператором. Требуется при использовании transformWithState оператора. Доступно в Databricks Runtime 16.4 LTS и более поздних версиях. |
readRegisteredTimers |
false |
true, false |
При trueчтении зарегистрированных таймеров, используемых оператором transformWithState . Применяется только к оператору transformWithState . Доступно в Databricks Runtime 16.4 LTS и более поздних версиях. |
flattenCollectionTypes |
true |
true, false |
Когда trueвозвращает записи, возвращаемые для переменных состояния карты и списка. Когда falseвозвращает записи в виде SQL Array Spark или Map. Применяется только к оператору transformWithState . Доступно в Databricks Runtime 16.4 LTS и более поздних версиях. |
Текст
При чтении текстовых файлов применяются следующие параметры.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
encoding |
UTF-8 |
java.nio.charset.Charset Имя |
Имя кодировки разделителя строк текстовых файлов. Содержимое файла не затрагивается этим параметром и считывается as-is. |
lineSep |
Нет, который охватывает \r, \r\n и \n |
Строка | Строка между двумя последовательными текстовыми записями. |
wholeText |
false |
true, false |
Следует ли считывать файл как одну запись. |
XML
При чтении XML-файлов применяются следующие параметры.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
rowTag |
None | Любая строка. | Тег строки в XML-файлах, который необходимо рассматривать как строку. В примере XML <book> <page><page>...<book>соответствующее значение имеет значение page. Это обязательный параметр. |
samplingRatio |
1.0 |
0.0 до 1.0 |
Определяет долю строк, используемых для вывода схемы. Встроенные функции XML игнорируют этот параметр. |
excludeAttribute |
false |
true, false |
Следует ли исключать атрибуты в элементах. |
mode |
None |
PERMISSIVE, , DROPMALFORMEDFAILFAST |
Режим работы с поврежденными записями во время синтаксического анализа.
|
inferSchema |
true |
true, false |
Если true, пытается определить соответствующий тип для каждого результирующего столбца DataFrame. Если false, все результирующие столбцы имеют тип string. Встроенные функции XML игнорируют этот параметр. |
columnNameOfCorruptRecord |
spark.sql.columnNameOfCorruptRecord |
Строка имени столбца | Позволяет переименовать новое поле, содержащее неправильно сформированную строку, созданную в режиме PERMISSIVE . |
attributePrefix |
None | Любая строка. | Префикс атрибутов для отличия атрибутов от элементов. Это будет префикс для имен полей. По умолчанию — _. Может быть пустым для чтения XML, но не для записи. Также относится к параметрам XML DataFrameWriter. |
valueTag |
_VALUE |
Любая строка. | Тег, используемый для символьных данных в элементах, которые также имеют атрибуты или дочерние элементы. Пользователь может указать поле valueTag в схеме, или оно будет добавлено автоматически во время определения схемы, если символьные данные присутствуют в элементах вместе с другими элементами или атрибутами. Также относится к параметрам XML DataFrameWriter. |
encoding |
UTF-8 |
java.nio.charset.Charset Имя |
Для чтения декодирует XML-файлы по заданному типу кодирования. Для записи задает кодировку (charset) сохраненных XML-файлов. Встроенные функции XML игнорируют этот параметр. Также относится к параметрам XML DataFrameWriter. |
ignoreSurroundingSpaces |
true |
true, false |
Следует ли пропускать пробелы, окружающие значения. Данные, содержащие только символьные пробелы, игнорируются. |
rowValidationXSDPath |
None | Строка пути к файлу | Путь к необязательному XSD-файлу, который используется для проверки XML для каждой строки по отдельности. Строки, которые не удается проверить, обрабатываются как ошибки синтаксического анализа. XSD не влияет на схему, указанную или выводимую. |
ignoreNamespace |
false |
true, false |
Если true, префиксы пространств имен для XML-элементов и атрибутов игнорируются. Теги <abc:author> и <def:author>, например, рассматриваются как если бы оба были просто <author>. Пространства имен нельзя игнорировать в элементе rowTag, можно игнорировать только читаемые дочерние элементы. Синтаксический анализ XML не учитывает пространство имен, даже если false. |
timestampFormat |
yyyy-MM-dd'T'HH:mm:ss[.SSS][XXX] |
Строка формата метки времени | Настраиваемая строка формата временной метки, которая соответствует формату шаблона даты и времени. Это относится к типу timestamp . Также относится к параметрам XML DataFrameWriter. |
timestampNTZFormat |
yyyy-MM-dd'T'HH:mm:ss[.SSS] |
Строка формата метки времени | Строка настраиваемого формата для метки времени без часового пояса, которая соответствует формату шаблона datetime. Это относится к типу TimestampNTZType. Также относится к параметрам XML DataFrameWriter. |
dateFormat |
yyyy-MM-dd |
Строка формата даты | Строка настраиваемого формата даты, которая соответствует шаблону формата datetime. Это относится к типу даты. Также относится к параметрам XML DataFrameWriter. |
locale |
en-US |
Тег языка IETF BCP 47 | Устанавливает локаль в виде тега языка в формате IETF BCP 47. Например, locale используется при анализе дат и меток времени. |
nullValue |
строка null |
Любая строка. | Задает строковое представление значения NULL. Когда это null, средство синтаксического анализа не записывает атрибуты и элементы для полей. Также относится к параметрам XML DataFrameWriter. |
readerCaseSensitive |
true |
true, false |
Указывает поведение чувствительности к регистру, когда включен параметр rescuedDataColumn. Если значение true, спасите столбцы данных, имена которых отличаются по регистру от схемы. Если значение false, считывает данные без учета регистра. |
rescuedDataColumn |
None | Строка имени столбца | Следует ли собирать все данные, которые нельзя проанализировать из-за несоответствия типов данных и несоответствия схемы (включая регистр столбцов) отдельному столбцу. Этот столбец включен по умолчанию при использовании Автозагрузчика. Для получения дополнительных сведений см. раздел "Что такое спасённый столбец данных?".
COPY INTO (устаревшая версия) не поддерживает спасательный столбец данных, так как невозможно вручную задать схему с помощью COPY INTO. Databricks рекомендует использовать Auto Loader для большинства сценариев загрузки. |
singleVariantColumn |
none |
Строка имени столбца | Указывает имя одного варианта столбца. Если этот параметр указан для чтения, синтаксический анализ всей XML-записи в один столбец Variant с заданным строковым значением параметра в качестве имени столбца. Если этот параметр указан для записи, напишите значение одного столбца Variant в XML-файлы. Также относится к параметрам XML DataFrameWriter. |
useLegacyXMLParser |
true |
true, false |
Следует ли использовать устаревший средство синтаксического анализа XML. Устаревший средство синтаксического анализа имеет менее строгие проверки для неправильно сформированного содержимого, но менее эффективной в памяти. Установите для false включения более строгого средства синтаксического анализа по умолчанию. |
wildcardColName |
xs_any |
Строка имени столбца | Имя столбца, используемое для записи XML-элементов, соответствующих элементу схемы подстановочного знака (xs:any). Нельзя использовать вместе с rescuedDataColumn. |
Параметры DataStreamReader
Используйте эти параметры для DataStreamReader.option() настройки потоковых операций чтения из таблиц Delta Lake и других источников на основе файлов.
Параметры формата файлов (JSON, CSV, Parquet и другие) см. в параметрах DataFrameReader.
Параметры автозагрузчика (cloudFiles.*) см. в разделе "Автозагрузчик".
Example
Следующий пример задает значение maxFilesPerTrigger10 для потока таблицы Delta Lake:
Python
df = spark.readStream.format("delta").option("maxFilesPerTrigger", 10).load("/path/to/delta-table")
Scala
val df = spark.readStream.format("delta").option("maxFilesPerTrigger", "10").load("/path/to/delta-table")
Общие
Следующие параметры применяются к таблицам Delta Lake и другим источникам потоковой передачи на основе файлов.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
cleanSource |
off |
off, , deletearchive |
Как обрабатывать исходные файлы после их обработки потоком.
off не принимает никаких действий.
delete окончательно удаляет исходный файл.
archive перемещает файл sourceArchiveDirв . Если задано значение archive, sourceArchiveDir также необходимо задать. Не применяется к потоковой передаче таблиц Delta Lake. |
fileNameOnly |
false |
true, false |
Следует ли определять уже обработанные файлы по имени файла, а не по полному пути. Если trueфайлы с разными путями с одинаковым именем файла обрабатываются как один и тот же файл и не обрабатываются повторно. Не применяется к потоковой передаче таблиц Delta Lake. |
latestFirst |
false |
true, false |
Следует ли обрабатывать последние измененные файлы в каждом микропакете. Полезно при обработке последних данных как можно быстрее. Если true и maxFilesPerTriggermaxBytesPerTrigger задано, maxFileAge игнорируется. Не применяется к потоковой передаче таблиц Delta Lake. |
maxBytesPerTrigger |
None | Положительные целые числа | Обратимое максимальное значение для объема данных, обрабатываемых для каждого микропакета. Пакет может обрабатывать больше предела, если наименьший входной блок превышает его. При использовании вместе с maxFilesPerTriggerданными микропакета обрабатываются до тех пор, пока не будет достигнуто любое ограничение.Для автозагрузчика используйте cloudFiles.maxBytesPerTrigger вместо этого. См. общие статьи. |
maxCachedFiles |
10000 |
Положительные целые числа или 0 |
Максимальное количество необработанных файлов для кэширования для последующих микропакетов. Установите для 0 отключения кэширования. Увеличьте это значение, если исходный каталог содержит множество новых файлов для каждого триггера. Не применяется к потоковой передаче таблиц Delta Lake. |
maxFileAge |
7d |
Строка длительности, 7d например или 4h |
Максимальный возраст файлов, которые считаются для обработки, относительно метки времени последнего измененного файла, а не текущего системного времени. Файлы старше этого порога игнорируются. Игнорируется, когда latestFirst задано true или maxFilesPerTriggermaxBytesPerTrigger задано. Не применяется к потоковой передаче таблиц Delta Lake. |
maxFilesPerTrigger |
1000 для Delta Lake и автозагрузчика. Нет максимального значения для других источников на основе файлов. |
Положительные целые числа | Верхняя граница для количества новых файлов, обработанных в каждом микропакете. При использовании вместе с maxBytesPerTriggerданными микропакета обрабатываются до тех пор, пока не будет достигнуто любое ограничение.Для автозагрузчика используйте cloudFiles.maxFilesPerTrigger вместо этого. См. общие статьи. |
sourceArchiveDir |
None | Строка пути | Путь к каталогу архива, если cleanSource задано значение archive. Исходные файлы перемещаются в этот путь после обработки, сохраняя относительную структуру каталога. Не применяется к потоковой передаче таблиц Delta Lake. |
Автозагрузчик
Используйте эти параметры с cloudFiles источником, чтобы настроить автозагрузчик для приема потоковой передачи из облачного хранилища. Параметры, относящиеся к cloudFiles источнику, имеют префикс, cloudFiles чтобы сохранить их в отдельном пространстве имен от других параметров источника структурированной потоковой передачи .
Общие
Следующие параметры применяются ко всем конфигурациям автозагрузчика.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
cloudFiles.allowOverwrites |
false |
true, false |
Разрешено ли изменение файла входного каталога для перезаписи существующих данных. Чтобы узнать о предостережениях при конфигурации, см. раздел Обрабатывает ли Автоматический загрузчик файл повторно при его добавлении или перезаписи?. |
cloudFiles.backfillInterval |
None | Строка длительности, 1 day например или 1 week |
Автозагрузчик может запускать асинхронные операции по заполнению данных с заданным интервалом. Дополнительные сведения см. в разделе "Активация обычных резервных копий с помощью cloudFiles.backfillInterval". Не используйте, если cloudFiles.useManagedFileEvents задано значение true. |
cloudFiles.cleanSource |
OFF |
OFF, , DELETEMOVE |
Следует ли автоматически удалять или перемещать обработанные файлы из входного каталога. Если задано значение OFF (по умолчанию), файлы не удаляются.В режиме DELETE автозагрузчик автоматически удаляет файлы через 30 дней после обработки. Для этого автозагрузчик должен иметь разрешения на запись в исходный каталог.При установке автозагрузчик MOVEавтоматически перемещает файлы в указанное расположение через cloudFiles.cleanSource.moveDestination 30 дней после обработки. Для этого автозагрузчик должен иметь разрешения на запись в исходный каталог, а также в расположение перемещения.Файл считается обработанным, если он имеет значение, отличное от NULL, в результате commit_time табличного значения cloud_files_state функции. См. табличную функцию cloud_files_state. 30-дневное дополнительное ожидание после обработки можно настроить с помощью cloudFiles.cleanSource.retentionDuration.Перед включением cloudFiles.cleanSourceознакомьтесь со следующими рекомендациями.
Доступно в Databricks Runtime 16.4 и более поздних версий. |
cloudFiles.cleanSource.retentionDuration |
30 days |
Строка CalendarInterval , например 14 days, 2 weeksили 1 month |
Время ожидания, прежде чем обработанные файлы становятся кандидатами для архивации.cleanSource Должно быть больше 7 дней для DELETE. Нет минимального ограничения для MOVE.Доступно в Databricks Runtime 16.4 и более поздних версий. |
cloudFiles.cleanSource.moveDestination |
None | Путь к облачному хранилищу или каталогу Unity | Путь к архиву для обработанных файлов, если для cloudFiles.cleanSource установлено значение MOVE. Это может быть путь к облачному хранилищу или путь к тому каталога Unity (например, /Volumes/my_catalog/my_schema/my_volume/archive/).Расположение перемещения должно:
Автозагрузчик должен иметь разрешения на запись в этот каталог. Доступно в Databricks Runtime 16.4 и более поздних версий. |
cloudFiles.format |
Нет (обязательный параметр) |
avro, binaryFile, csvjsonorcparquettextxml |
Формат файла данных в исходном пути. Допустимые значения:
|
cloudFiles.includeExistingFiles |
true |
true, false |
Следует ли включать существующие файлы во входной путь обработки потоковой передачи или обрабатывать только новые файлы, поступающие после первоначальной настройки. Этот параметр оценивается только при первом запуске потока. Изменение этого параметра после перезапуска потока не даст результата. |
cloudFiles.inferColumnTypes |
false |
true, false |
Следует ли выводить точные типы столбцов при использовании вывода схемы. По умолчанию столбцы выводятся как строки при выводе наборов данных JSON и CSV. Дополнительные сведения см. в разделе Вывод схемы. |
cloudFiles.maxBytesPerTrigger |
None | Строка байтов, например 10g |
Максимальное число новых байтов, которое может обрабатываться в каждом триггере. Это мягкое ограничение. Если у вас есть файлы размером 3 ГБ, Azure Databricks обрабатывает 12 ГБ в микропакете. Отдельный файл никогда не разбивается по микропакетам; он всегда обрабатывается в полном объеме в пределах одного, даже если его размер превышает этот предел. При совместном использовании с cloudFiles.maxFilesPerTrigger Azure Databricks потребляет до нижнего предела cloudFiles.maxFilesPerTrigger или cloudFiles.maxBytesPerTrigger, в зависимости от того, что будет достигнуто раньше. Этот параметр не действует при использовании с Trigger.Once() (Trigger.Once() не рекомендуется).В Databricks Runtime 18.0 и более поздних версий этот параметр настраивается динамически и не требуется устанавливать вручную. |
cloudFiles.maxFileAge |
None | Строка длительности | Длительность отслеживания события файла для обнаружения дубликатов. В Databricks не рекомендуется настраивать этот параметр, если вы не загружаете данные в размере миллионов файлов в час. Дополнительные сведения см. в разделе "Отслеживание событий файлов ". Настройка cloudFiles.maxFileAge слишком агрессивно может вызвать проблемы с качеством данных, такие как дублирование загрузки или пропущенные файлы. Таким образом, Databricks рекомендует установить консервативное значение для cloudFiles.maxFileAge, например, 90 дней, что аналогично рекомендациям сопоставимых решений для приема данных. |
cloudFiles.maxFilesPerTrigger |
1000 |
Положительные целые числа | Максимальное число новых файлов, которое должно быть обработано в каждом триггере. При совместном использовании с cloudFiles.maxBytesPerTrigger Azure Databricks потребляет до нижнего предела cloudFiles.maxFilesPerTrigger или cloudFiles.maxBytesPerTrigger, в зависимости от того, что будет достигнуто раньше. Этот параметр не действует при использовании с Trigger.Once() (не рекомендуется).В Databricks Runtime 18.0 и более поздних версий этот параметр настраивается динамически и не требуется устанавливать вручную. |
cloudFiles.partitionColumns |
None | Разделенный запятыми список имен столбцов | Разделенный запятыми список столбцов секций в стиле Hive, которые вы хотите вывести из структуры каталогов файлов. Столбцы секционирования в стиле Hive — это пары "ключ-значение", объединенные знаком равенства, например <base-path>/a=x/b=1/c=y/file.format. В этом примере столбцы секционирования представляют собой a, b и c. По умолчанию эти столбцы автоматически добавляются в схему при использовании вывода схемы и указывают <base-path> данные для загрузки данных. Если указать схему, автозагрузчик ожидает, что эти столбцы будут включены в схему. Если вы не хотите, чтобы эти столбцы были включены в схему, можно указать "", чтобы игнорировать их. Кроме того, этот параметр можно использовать, если требуется, чтобы столбцы выводили путь к файлу в сложных структурах каталогов, как показано в примере ниже.<base-path>/year=2022/week=1/file1.csv<base-path>/year=2022/month=2/day=3/file2.csv<base-path>/year=2022/month=2/day=4/file3.csvУказание cloudFiles.partitionColumns в качестве year,month,day возвращает year=2022 значения для file1.csv, но столбцы month и daynull.month и day правильно анализируются для file2.csv и file3.csv. |
cloudFiles.schemaEvolutionMode |
addNewColumns Если схема не указана, none в противном случае |
addNewColumns, , nonerescuefailOnNewColumns |
Режим для развития схемы по мере обнаружения в данных новых столбцов. По умолчанию столбцы выводятся как строки при выводе наборов данных JSON. Дополнительные сведения см. в разделе Развитие схемы. |
cloudFiles.schemaHints |
None | Строка схемы | Сведения о схеме, указанные для автозагрузчика во время вывода схемы. Дополнительные сведения см. в разделе Указания для схемы. |
cloudFiles.schemaLocation |
Нет (требуется для вывода схемы) | Строка пути | Расположение для хранения выводимой схемы и последующих изменений. Дополнительные сведения см. в разделе Вывод схемы. |
cloudFiles.useStrictGlobber |
false |
true, false |
Следует ли использовать строгий режим глоббинга, соответствующий поведению глоббинга других источников файлов по умолчанию в Apache Spark. Дополнительные сведения см. в общих шаблонах загрузки данных . Доступно в Databricks Runtime 12.2 LTS и более поздних версиях. |
cloudFiles.validateOptions |
true |
true, false |
Следует ли проверять параметры Автозагрузчика и возвращать ошибку для неизвестных или несогласованных параметров. |
Список каталогов
Следующий параметр применяется при использовании режима перечисления каталогов.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
cloudFiles.useIncrementalListing (не рекомендуется) |
auto в Databricks Runtime 17.2 и более поздних версиях в false Databricks Runtime 17.3 и более поздних версиях |
auto, , truefalse |
Эта функция является устаревшей. Databricks рекомендует использовать режим уведомлений файлов с событиями файлов вместо cloudFiles.useIncrementalListing.Следует ли использовать добавочный список вместо полного списка в режиме отображения содержимого каталога. По умолчанию автозагрузчик делает все возможное, чтобы автоматически определить, применим ли данный каталог к добавочному перечислению. Вы можете явно использовать добавочный список или использовать полный список файлов каталогов, задав для этого параметра значение true или false соответственно.Неправильное включение инкрементального перечисления в нелексически упорядоченном каталоге мешает автозагрузчику обнаруживать новые файлы. Работает с Azure Data Lake Storage (), S3 ( abfss://s3://) и GCS (gs://).Доступно в Databricks Runtime 9.1 LTS и более поздних версиях. |
Уведомление о файле
Сведения о настройке режима уведомлений файлов, включая необходимые облачные разрешения, инструкции по настройке и методы проверки подлинности, см. в разделе "Настройка потоков автозагрузчика" в режиме уведомлений о файлах.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
cloudFiles.fetchParallelism |
1 |
Положительные целые числа | Число потоков, используемых для извлечения сообщений из системы управления очередями. Не используйте, если cloudFiles.useManagedFileEvents задано значение true. |
cloudFiles.pathRewrites |
None | Строка карты JSON | Требуется только в том случае, если вы указываете queueUrl , что получает уведомления о файлах из нескольких контейнеров S3 и хотите использовать точки подключения, настроенные для доступа к данным в этих контейнерах. Используйте этот параметр, чтобы заменить префикс bucket/key пути на точку монтирования. Перезаписывать можно только префиксы. Например, для конфигурации {"<databricks-mounted-bucket>/path": "dbfs:/mnt/data-warehouse"}путь s3://<databricks-mounted-bucket>/path/2017/08/fileA.json перезаписан в dbfs:/mnt/data-warehouse/2017/08/fileA.json.Не используйте, если cloudFiles.useManagedFileEvents задано значение true. |
cloudFiles.resourceTag |
None | Строки тегов key-value | Ряд пар тегов "ключ — значение", которые помогают связать и идентифицировать связанные ресурсы, например:cloudFiles.option("cloudFiles.resourceTag.myFirstKey", "myFirstValue") .option("cloudFiles.resourceTag.mySecondKey", "mySecondValue")Не используйте, если cloudFiles.useManagedFileEvents задано значение true. Вместо этого задайте теги ресурсов с помощью консоли поставщика облачных служб.Дополнительные сведения см. в разделе тегов ресурсов поставщика облачных служб. |
cloudFiles.useManagedFileEvents |
false |
true, false |
При установке true автозагрузчик использует службу событий файлов для обнаружения файлов во внешнем расположении. Этот параметр можно использовать только в том случае, если путь загрузки данных находится во внешней директории, где включены события файлов. См. раздел "Использование режима уведомлений файлов" с событиями файлов.События файлов обеспечивают производительность на уровне уведомлений при обнаружении файлов, так как автозагрузчик может обнаруживать новые файлы после последнего запуска. В отличие от списка каталогов, этот процесс не должен перечислять все файлы в каталоге. Существуют некоторые ситуации, когда автозагрузчик использует список каталогов, даже если включен параметр событий файла:
См. раздел "Когда автозагрузчик с событиями файлов использует список каталогов?", чтобы получить полный список ситуаций, когда автозагрузчик использует список каталогов с этим параметром. Доступно в Databricks Runtime 14.3 LTS и более поздних версиях. |
cloudFiles.listOnStart |
false |
true, false |
Если задано значение true, автозагрузчик выполняет полный список каталогов при запуске потока, а не начиная с маркера продолжения в контрольной точке. Используйте этот параметр для восстановления после ошибок, таких как CF_MANAGED_FILE_EVENTS_INVALID_CONTINUATION_TOKEN. Узнайте, как восстановиться после CF_MANAGED_FILE_EVENTS_INVALID_CONTINUATION_TOKEN ошибки? |
cloudFiles.useNotifications |
false |
true, false |
Следует ли использовать режим уведомлений о файлах для определения наличия новых файлов. Если значение равно false, используется режим листинга каталогов. См. Сравнение режимов обнаружения файлов автозагрузчика.Не используйте, если cloudFiles.useManagedFileEvents задано значение true. |
Теги ресурсов поставщика облачных служб
Автозагрузчик добавляет следующие пары тегов "ключ-значение" по умолчанию на основе наилучших усилий:
-
vendor:Databricks -
path: расположение, из которого загружаются данные. Недоступно в GCP из-за ограничений по маркировке. -
checkpointLocation: расположение контрольной точки потока. Недоступно в GCP из-за ограничений по маркировке. -
streamId: глобальный уникальный идентификатор для потока.
Databricks резервирует эти имена ключей, и вы не можете перезаписать их значения.
Дополнительные сведения об Azure см. в разделе Именование очередей и метаданных и охват properties.labels в подписках на события. Автозагрузчик сохраняет эти пары тегов "ключ — значение" в формате JSON в виде меток.
Облачные решения
Автозагрузчик имеет параметры настройки облачной инфраструктуры для режима уведомлений о файлах. Необходимые облачные разрешения и инструкции по настройке см. в разделе "Настройка потоков автозагрузчика" в режиме уведомлений о файлах.
Azure
При указании необходимо указать значения для всех следующих параметров, и cloudFiles.useNotifications = true вы хотите, чтобы автозагрузчик настраивал службы уведомлений для вас:
Если учетные данные службы Databricks недоступны, можно указать следующие параметры проверки подлинности.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
cloudFiles.clientId |
None | Любая строка. | Идентификатор клиента или приложения сервисного принципала. |
cloudFiles.clientSecret |
None | Любая строка. | Секрет клиента субъекта-службы. |
cloudFiles.connectionString |
None | Строка подключения | Строка подключения для учетной записи хранения, которая может быть основана на ключе доступа к учетной записи или токене совместного доступа (SAS). |
cloudFiles.tenantId |
None | Любая строка. | Идентификатор клиента Azure, в котором создается субъект-служба. |
Укажите следующий параметр, только если задано cloudFiles.useNotifications = true и требуется, чтобы автозагрузчик использовал существующую очередь:
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
cloudFiles.queueName |
None | Любая строка. | Имя очереди Azure. Если задано, источник облачных файлов напрямую использует события из этой очереди, а не настраивает собственные службы Сетка событий Azure и хранилища очередей. В этом случае для databricks.serviceCredential или cloudFiles.connectionString требуются только права на чтение в очереди. |
Delta Lake
Следующие параметры применяются при чтении из таблицы Delta Lake с помощью spark.readStream.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
allowSourceColumnDrop |
None | Номер версии или always |
Задайте номер версии таблицы Delta или always позволить потоку продолжаться после удаления столбцов из схемы исходной таблицы. Если задано значение номера версии, все изменения схемы изменяются до этой версии. Требует использования schemaTrackingLocation. См. как переименовывать и удалять столбцы с использованием сопоставления столбцов в Delta Lake. |
allowSourceColumnRename |
None | Номер версии или always |
Задайте номер версии таблицы Delta или always разрешить потоку продолжаться после переименования столбцов в исходной таблице. Если задано значение номера версии, все изменения схемы изменяются до этой версии. Требует использования schemaTrackingLocation. См. как переименовывать и удалять столбцы с использованием сопоставления столбцов в Delta Lake. |
allowSourceColumnTypeChange |
None | Номер версии или always |
Задайте номер версии таблицы Delta или always разрешить потоку продолжаться после изменения типов столбцов в исходной таблице. Если задано значение номера версии, все изменения схемы изменяются до этой версии. Требует использования schemaTrackingLocation. См. Расширение типов. |
excludeRegex |
None | Строка regex Java | Шаблон регулярного выражения. Файлы, пути которых соответствуют шаблону, исключаются из потоковой передачи. Полезно для фильтрации файлов, которые не соответствуют ожидаемому соглашению об именовании. |
failOnDataLoss |
true |
true, false |
Если исходные данные были удалены из-за хранения журналов (logRetentionDuration). Установите для false пропуска отсутствующих данных и продолжения обработки. См. Настройка сохранения данных для запросов по временному перемещению. |
ignoreChanges (не рекомендуется) |
false |
true, false |
Доступно в Databricks Runtime 11.3 LTS и ниже. Повторно отправляет файлы данных после операций изменения, таких как UPDATE, MERGE INTOили DELETEOVERWRITE. Без изменений строки могут создаваться вместе с новыми строками, поэтому подчиненные потребители должны обрабатывать дубликаты. Удаления не распространяются вниз по потоку. Заменено skipChangeCommits в Databricks Runtime 12.2 LTS и выше. |
ignoreDeletes (не рекомендуется) |
false |
true, false |
Игнорирует транзакции, которые удаляют данные по границам секции (только полное удаление секций). Не обрабатывает удаления, обновления и другие изменения, не относящиеся к секционированиям. Вместо этого используйте skipChangeCommits. |
readChangeFeed или readChangeData |
false |
true, false |
Следует ли включить чтение веб-канала измененных данных для потокового запроса. При включении поток выдает изменения на уровне строк (вставки, обновления и удаления) с дополнительными столбцами метаданных. См. Использование канала передачи данных об изменениях в Azure Databricks. |
schemaTrackingLocation |
None | Строка пути | Путь к каталогу, в котором Delta Lake отслеживает изменения схемы для потокового чтения. Требуется при потоковой передаче из таблиц с включенным сопоставлением столбцов и использованием allowSourceColumn* параметров для обработки эволюции схемы. Должен находиться в checkpointLocation потоковом запросе. См. как переименовывать и удалять столбцы с использованием сопоставления столбцов в Delta Lake. |
skipChangeCommits |
false |
true, false |
Игнорирует транзакции, которые удаляют или изменяют существующие записи и процессы только добавляются. Databricks рекомендует этот параметр для большинства рабочих нагрузок, которые не используют веб-каналы измененных данных. Доступно в Databricks Runtime 12.2 LTS и более поздних версиях. См. раздел "Пропустить фиксации изменений вышестоящего потока" с skipChangeCommitsпомощью . |
startingTimestamp |
Последняя версия доступна | Строка метки времени, 2019-01-01T00:00:00.000Z например или строка даты, например 2019-01-01 |
Метка времени, из которой начинается чтение. Поток считывает все изменения таблицы, зафиксированные в указанной метке времени или после нее. Если метка времени предшествует всем доступным фиксациям таблицы, поток начинается с самой ранней доступной фиксации. Нельзя использовать вместе с startingVersion. Игнорируется, если контрольная точка потоковой передачи уже существует. |
startingVersion |
Последняя версия доступна | Положительное целое число или 0latest |
Разностная версия таблицы, из которой начинается чтение. Поток считывает все изменения, зафиксированные в указанной версии или после нее. Укажите latest , чтобы начать только последние изменения. Нельзя использовать вместе с startingTimestamp. Игнорируется, если контрольная точка потоковой передачи уже существует. См. статью "Работа с журналом таблиц". |
withEventTimeOrder |
false |
true, false |
Делит начальный моментальный снимок таблицы на контейнеры времени событий, чтобы предотвратить неправильно помеченные записи как поздние события и удалены в запросах с отслеживанием состояния с подложками. Невозможно изменить после начала начальной обработки моментальных снимков без удаления контрольной точки. Доступно в Databricks Runtime 11.3 LTS и более поздних версиях. См. также Обработка начального моментального снимка без потери данных. |
Кафка
Используйте следующие параметры с помощью следующих spark.readStream.format("kafka") вариантов:spark.read.format("kafka")
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
assign |
None | Строка JSON, например {"topicA":[0,1],"topicB":[2,4]} |
Определенные секции для использования. Необходимо указать именно один из subscribesubscribePatternпараметров или assign параметров. |
failOnDataLoss |
true |
true, false |
Не удается выполнить запрос, если данные могли быть потеряны, например из-за удаленных разделов или усечения смещения. Установите для false пропуска отсутствующих данных и продолжения.Databricks оценивает консервативно, могут ли данные быть потеряны. Однако это может привести к ложным тревогам. |
fetchoffset.numretries |
3 |
Положительные целые числа или 0 |
Количество повторных попыток при получении смещения Kafka завершается сбоем. |
fetchoffset.retryintervalms |
1000 |
Положительные целые числа или 0 |
Интервал в миллисекундах между повторными попытками получения смещения. |
groupIdPrefix |
spark-kafka-source (потоковая передача), spark-kafka-relation (пакетная версия) |
Любая строка. | Настраиваемый префикс для автоматического создания идентификатора группы потребителей Kafka. Если kafka.group.id этот параметр задан явным образом, соединитель игнорирует этот параметр. |
kafka.group.id |
None | Любая строка. | Идентификатор группы потребителей Kafka, используемый при чтении. Используйте с осторожностью: запросы, использующие один и тот же идентификатор группы, вмешиваются друг в друга и могут читать только частичные данные. Это может произойти при выполнении параллельных рабочих нагрузок пакетной и потоковой передачи или при быстром перезапуске запросов. Если задано, groupIdPrefix игнорируется. Чтобы свести к минимуму проблемы, задайте конфигурацию session.timeout.ms потребителя Kafka небольшим значением. |
includeHeaders |
false |
true, false |
Следует ли включать заголовки сообщений Kafka в качестве столбца в выходные данные. |
kafkaconsumer.polltimeoutms |
None | Положительные целые числа | Время ожидания в миллисекундах для вызова потребителя poll() Kafka. |
kafka.bootstrap.servers |
None | Разделенный запятыми список host:port строк |
Разделенный запятыми список адресов host:port для брокеров Kafka. Задает свойство клиента bootstrap.servers Kafka.Если данные из Kafka отсутствуют, проверьте этот список адресов брокера на наличие неправильных адресов. Если список адресов брокера неверный, ошибки могут не возникнуть. Клиенты Kafka предполагают, что брокеры будут доступны в конечном итоге и повторяться навсегда, когда они получают сетевые ошибки. |
maxRecordsPerPartition |
None | Положительные целые числа | Максимальное количество записей для каждой секции Spark. При установке соединитель разделяет секции Kafka, чтобы каждая секция Spark считывала не больше всего этих записей. Этот параметр также можно использовать с minPartitions. При задании обоих параметров Spark использует любой параметр, который приводит к дополнительным секциям. |
minPartitions |
None | Положительные целые числа | Минимальное количество секций Spark для чтения из Kafka. При установке соединитель разбивает большие секции Kafka, чтобы увеличить параллелизм. Если не задано, Spark создает одну секцию для каждой секции раздела Kafka. Полезно для обработки отклонений данных или пиковых нагрузок. Этот параметр повторно инициализирует потребителей Kafka для каждого триггера, что может повлиять на производительность с помощью SSL. |
startingOffsets |
latest (потоковая передача), earliest (пакетная версия) |
earliest, latestили строка смещения JSON |
Смещение, из которое начинается запрос. В строке -1 JSON используется последнее смещение.
-2 является самым ранним смещением. Например: {"topicA":{"0":23,"1":-2}}.Для потоковых запросов этот параметр применяется только при запуске нового запроса. Возобновленные запросы всегда используют контрольную точку. Во время запроса новые секции начинают читаться с самого раннего смещения. Для пакетных latest запросов запрещено. |
startingOffsetsByTimestamp |
None | Строка метки времени JSON, например {"topicA":{"0":1000,"1":2000}} |
Список начальных смещения для каждой секции, указанный как метки времени в миллисекундах. Если для метки времени не существует смещения, поведение запроса определяется startingOffsetsByTimestampStrategy.Для потоковых запросов этот параметр применяется только при запуске нового запроса. Возобновленные запросы всегда используют контрольную точку. Во время запроса новые секции начинают читаться с самого раннего смещения. |
startingOffsetsByTimestampStrategy |
error |
error, latest |
Стратегия, используемая при отсутствии смещения для метки времени, указанной в startingOffsetsByTimestamp или startingTimestamp.
error вызывает исключение.
latest использует последнее доступное смещение. |
startingTimestamp |
None | Положительные целые числа или 0 |
Глобальная метка времени запуска в миллисекундах, которая применяется ко всем секциям. Если для метки времени не существует смещения, поведение управляется startingOffsetsByTimestampStrategy. |
subscribe |
None | Разделенный запятыми список имен разделов | Разделы, на которые нужно подписаться. Необходимо указать именно один из subscribesubscribePatternпараметров или assign параметров. |
subscribePattern |
None | Строка regex Java | Шаблон, используемый для подписки на разделы. Необходимо указать именно один из subscribesubscribePatternпараметров или assign параметров. Например: topic.*. |
Следующие параметры применяются только к потоковым чтениям со spark.readStream.format("kafka")следующими параметрами:
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
bytesEstimateWindowLength |
300s |
Строки длительности, 10m например или 600s |
Период времени, используемый для оценки оставшихся байтов для estimatedTotalBytesBehindLatest метрики. Ознакомьтесь с получением метрик Kafka. |
maxOffsetsPerTrigger |
None | Положительные целые числа | Максимальное количество смещения для обработки интервала триггера. Смещения распределяются пропорционально между секциями раздела. |
maxTriggerDelay |
15m |
Строки длительности, 10m например или 600s |
Максимальное время ожидания minOffsetsPerTrigger до активации. |
minOffsetsPerTrigger |
None | Положительные целые числа | Минимальное количество смещения для накапливания перед активацией микропакета. По maxTriggerDelay достижении микропакет выполняется независимо от того, |
Параметры смещения, которые применяются только к пакетным считываниям, spark.read.format("kafka")см. в параметрах Kafka DataFrameReader.
Проверка подлинности
Databricks рекомендует использовать учетные данные службы каталога Unity для проверки подлинности в облачных службах Kafka (AWS MSK, Центры событий Azure или Google Cloud Managed Kafka).
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
databricks.serviceCredential |
None | Любая строка. | Имя учетных данных службы каталога Unity для проверки подлинности в облачных службах Kafka. Доступно в Databricks Runtime 16.1 и более поздних версиях. |
databricks.serviceCredential.scope |
None | Любая строка. | Область OAuth для учетных данных службы. Задайте это только в том случае, если Azure Databricks не может автоматически выводить область для службы Kafka. |
Если учетные данные службы недоступны, используйте параметры SASL/SSL (передаваемые в качестве kafka.* свойств). При использовании учетных данных службы не требуется указывать kafka.sasl.mechanismилиkafka.sasl.jaas.configkafka.security.protocol.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
kafka.security.protocol |
None | Строка протокола безопасности, например SASL_SSL, SSLPLAINTEXT |
Протокол безопасности для обмена данными брокера. |
kafka.sasl.mechanism |
None | Строка механизма SASL, например PLAIN, SCRAM-SHA-256, SCRAM-SHA-512OAUTHBEARERAWS_MSK_IAM |
Механизм SASL. |
kafka.sasl.jaas.config |
None | Строка конфигурации JAAS | Строка конфигурации входа JAAS. |
kafka.sasl.login.callback.handler.class |
None | Полное имя класса | Полное имя класса обработчика обратного вызова для проверки подлинности SASL. |
kafka.sasl.client.callback.handler.class |
None | Полное имя класса | Полное имя класса обработчика обратного вызова клиента для проверки подлинности SASL. |
kafka.ssl.truststore.location |
None | Строка пути к файлу | Путь к файлу хранилища доверия SSL. |
kafka.ssl.truststore.password |
None | Любая строка. | Пароль для файла хранилища доверия SSL. |
kafka.ssl.keystore.location |
None | Строка пути к файлу | Путь к файлу хранилища ключей SSL. |
kafka.ssl.keystore.password |
None | Любая строка. | Пароль для файла хранилища ключей SSL. |
Полные инструкции по настройке проверки подлинности см. в разделе "Проверка подлинности".
Pub/Sub
Используйте эти параметры, spark.readStream.format("pubsub") чтобы подписаться на Google Pub/Sub. Параметры subscriptionIdи topicIdprojectId обязательные.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
subscriptionId |
None | Любая строка. | Обязательно. Идентификатор подписки Pub/Sub. Соединитель создает подписку, если она не существует. |
topicId |
None | Любая строка. | Обязательно. Идентификатор раздела Pub/Sub. |
projectId |
None | Любая строка. | Обязательно. Идентификатор проекта Google Cloud. |
numFetchPartitions |
Половина числа исполнителей, доступных при инициализации потока | Положительные целые числа | Количество параллельных задач Spark, которые извлекает строки из подписки. |
maxBytesPerTrigger |
None | Положительные целые числа | Обратимое ограничение на количество байтов для обработки на микропакет. |
maxRecordsPerFetch |
1000 |
Положительные целые числа | Количество строк для получения каждой задачи перед обработкой. |
maxFetchPeriod |
10s |
Строка длительности, 1s например или 1m |
Длительность времени для каждой задачи, извлекаемой перед обработкой строк. Azure Databricks рекомендует использовать значение по умолчанию. |
deleteSubscriptionOnStreamStop |
false |
true, false |
Когда trueподписка subscriptionIdудаляется при завершении потокового запроса. |
serviceCredential |
None | Любая строка. | Имя учетных данных службы Azure Databricks для проверки подлинности в Pub/Sub. Доступно в Databricks Runtime 16.1 и более поздних версиях. |
clientEmail |
None | Строка адреса электронной почты | Адрес электронной почты учетной записи службы Google. Требуется, если учетные данные службы не используются. |
clientId |
None | Любая строка. | Идентификатор клиента учетной записи службы Google. Требуется, если учетные данные службы не используются. |
privateKey |
None | Строка закрытого ключа | Закрытый ключ для учетной записи службы Google. Требуется, если учетные данные службы не используются. |
privateKeyId |
None | Любая строка. | Идентификатор закрытого ключа для учетной записи службы Google. Требуется, если учетные данные службы не используются. |
Дополнительные сведения о pub/Sub см. в разделе "Подписка на Google Pub/Sub".
Пульсар
Используйте эти параметры для spark.readStream.format("pulsar") потоковой передачи из Apache Pulsar. Доступно в Databricks Runtime 14.1 и более поздних версиях.
Необходимы следующие параметры. Необходимо указать именно один из topic, topicsили topicsPattern.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
service.url |
None | Строка URL-адреса службы Pulsar | Например, serviceURLPulsar для службы Pulsarpulsar://broker.example.com:6650. |
topic |
None | Любая строка. | Имя одного раздела для использования. |
topics |
None | Разделенный запятыми список имен разделов | Разделенный запятыми список имен разделов для использования. |
topicsPattern |
None | Строка regex Java | Строка Java regex для сопоставления имен разделов. |
Также поддерживаются следующие параметры:
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
admin.url |
None | Строка URL-адреса | HTTP-адрес службы администрирования Pulsar. Обязательный параметр при maxBytesPerTrigger установке. |
allowDifferentTopicSchemas |
false |
true, false |
Если считываются несколько разделов с разными схемами, используйте этот параметр, чтобы отключить автоматическую десериализацию значений раздела на основе схем. Возвращаются только необработанные значения, если это true. |
failOnDataLoss |
true |
true, false |
Происходит ли сбой запроса при потере данных. Например, потеря данных может возникать при удалении тем или истечении срока действия сообщений из-за политики хранения. |
maxBytesPerTrigger |
None | Положительные целые числа | Обратимое ограничение на количество байтов для обработки на микропакет. Требует использования admin.url. |
pollTimeoutMs |
120000 |
Положительные целые числа | Время ожидания для чтения сообщений из Pulsar в миллисекундах. |
predefinedSubscription |
None | Любая строка. | Предопределенное имя подписки, используемое соединителем для отслеживания хода выполнения приложения Spark. |
startingOffsets |
latest |
latest, earliestили строка смещения JSON |
Откуда начать чтение. |
subscriptionPrefix |
None | Любая строка. | Префикс, используемый соединителем для создания случайной подписки для отслеживания хода выполнения приложения Spark. |
waitingForNonExistedTopic |
false |
true, false |
Ожидается ли соединитель до создания нужных разделов. |
Вы можете указать дополнительные конфигурации клиента, администратора и читателя Pulsar, используя следующие шаблоны параметров:
| Рисунок | Варианты конфигурации |
|---|---|
pulsar.admin.* |
Конфигурация администратора Pulsar |
pulsar.client.* |
Конфигурация клиента Pulsar, включая такие параметры проверки подлинности, как pulsar.client.authPluginClassName и pulsar.client.authParams. |
pulsar.reader.* |
Конфигурация средства чтения Pulsar |
Дополнительные сведения о параметрах проверки подлинности клиента и администратора Pulsar см. в разделе "Проверка подлинности".
Проверка подлинности
Azure Databricks поддерживает аутентификацию для Pulsar с использованием хранилища доверенных сертификатов и хранилища ключей. Azure Databricks рекомендует использовать секреты для хранения сведений о проверке подлинности. Дополнительные сведения см. в разделе Управление секретами.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
pulsar.client.authPluginClassName |
None | Полное имя класса | Полное имя класса подключаемого модуля проверки подлинности. Например: org.apache.pulsar.client.impl.auth.AuthenticationTls. |
pulsar.client.authParams |
None | Строка учетных данных | Учетные данные проверки подлинности, передаваемые в подключаемый модуль проверки подлинности в виде строки. Например: tlsCertFile:/path/to/my-role.cert.pem,tlsKeyFile:/path/to/my-role.key-pk8.pem. |
pulsar.client.useKeyStoreTls |
false |
true, false |
При trueвключении конфигурации TLS на основе KeyStore вместо файлов формата PEM. |
pulsar.client.tlsTrustStoreType |
None | Любая строка. | Формат файла хранилища доверия TLS. Например: JKS. |
pulsar.client.tlsTrustStorePath |
None | Строка пути к файлу | Путь к файлу хранилища доверия TLS, содержанию доверенных сертификатов ЦС. Обязательный, если pulsar.client.useKeyStoreTls имеет значение true. |
pulsar.client.tlsTrustStorePassword |
None | Любая строка. | Пароль для файла хранилища доверия TLS. |
Если в потоке используется, PulsarAdminможно также задать следующие параметры:
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
pulsar.admin.authPluginClassName |
None | Полное имя класса | Полное имя класса подключаемого модуля проверки подлинности для клиента администратора Pulsar. |
pulsar.admin.authParams |
None | Строка учетных данных | Учетные данные проверки подлинности для подключаемого модуля проверки подлинности клиента администратора Pulsar. |
pulsar.admin.useTls |
None |
true, false |
Следует ли использовать TLS для подключения клиента администратора Pulsar. |
pulsar.admin.tlsAllowInsecureConnection |
None |
true, false |
Следует ли разрешать небезопасные подключения TLS для клиента администратора Pulsar. |
pulsar.admin.tlsTrustCertsFilePath |
None | Строка пути к файлу | Путь к доверенному файлу сертификата TLS для клиента администратора Pulsar. |
pulsar.admin.useKeyStoreTls |
None |
true, false |
Следует ли использовать TLS на основе KeyStore для клиента администратора Pulsar. |
pulsar.admin.tlsTrustStoreType |
None | Любая строка. | Формат хранилища доверия TLS для клиента администратора Pulsar. Например: JKS. |
pulsar.admin.tlsTrustStorePath |
None | Строка пути к файлу | Путь к файлу хранилища доверия TLS для клиента администратора Pulsar. Обязательный, если pulsar.admin.useKeyStoreTls имеет значение true. |
pulsar.admin.tlsTrustStorePassword |
None | Любая строка. | Пароль для хранилища доверия клиента администратора Pulsar. |
Примеры проверки подлинности см. в разделе "Проверка подлинности в Pulsar".
Параметры DataFrameWriter
Используйте эти параметры с DataFrameWriter.option() и DataFrameWriterV2.option() для управления Azure Databricks записи данных.
Example
В следующем примере задано mergeSchema значение True для записи таблицы Delta Lake:
Python
df.write.format("delta").option("mergeSchema", True).saveAsTable("my_table")
Scala
df.write.format("delta").option("mergeSchema", "true").saveAsTable("my_table")
Avro
При написании файлов Avro применяются следующие параметры.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
avroSchema |
None | Строка схемы JSON | Полная схема Avro в виде строки JSON. Используйте этот параметр, чтобы преобразовать типы SQL Spark в определенные типы Avro. Применяется к файлам Avro чтения и записи. |
avroSchemaUrl |
None | Строка URL-адреса | URL-адрес, указывающий на файл схемы Avro. Используйте вместо того, avroSchema когда схема хранится внешне. Взаимоисключающ с avroSchema. Применяется к файлам Avro чтения и записи. |
compression |
snappy |
uncompressed, , deflatesnappy (default)bzip2xz,zstandard |
Кодек сжатия, используемый при записи. Применяется к файлам Avro чтения и записи. |
recordName |
topLevelRecord |
Любая строка. | Имя записи верхнего уровня в выходной схеме Avro. Применяется к файлам Avro чтения и записи. |
positionalFieldMatching |
false |
true, false |
Следует ли сопоставлять столбцы между схемой Spark и схемой Avro по позиции поля вместо имени. Применяется к файлам Avro чтения и записи. |
recordNamespace |
Пустая строка | Любая строка. | Пространство имен для записи верхнего уровня в выходной схеме Avro. Применяется к файлам Avro чтения и записи. |
Delta Lake и Apache Iceberg
Следующие параметры применяются при написании таблиц Delta Lake и Apache Iceberg.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
clusterByAuto |
false |
true, false |
Следует ли включить автоматическую кластеризацию жидкости, где Azure Databricks выбирает столбцы кластеризации на основе шаблонов запросов. Допустимо только с mode("overwrite"). Нельзя использовать с append режимом. Доступно в Databricks Runtime 16.4 и более поздних версий. Применяется к использованию отказоустойчивой кластеризации для таблиц. |
mergeSchema |
None |
true, false |
Следует ли включить эволюцию схемы для операции записи. Новые столбцы в исходном кадре данных добавляются в целевую схему таблицы. Применяется к добавкам пакетной и потоковой передачи. Применяется к схемам таблиц обновления с эволюцией схемы. |
overwriteSchema |
None |
true, false |
Следует ли заменить схему таблицы и секционирование при перезаписи. Требуется mode("overwrite") без replaceWhere. Нельзя использовать с partitionOverwriteMode. Применяется к схемам таблиц обновления с эволюцией схемы. |
partitionOverwriteMode |
None |
static, dynamic |
Режим перезаписи секции. Установите для этого значение, чтобы dynamic перезаписать только разделы, содержащие новые данные, оставляя все остальные секции без изменений. Устаревший режим, не поддерживаемый для бессерверных вычислений или Databricks SQL. Применяется к выборочному перезаписи данных с помощью Delta Lake. |
replaceOn |
None | Логическое выражение | Логическое выражение, которое соответствует строкам в целевой таблице для замены строк из исходного запроса. Может ссылаться на столбцы из целевой таблицы и исходного запроса. Строки в целевом объекте, который соответствует исходной строке, удаляются и заменяются. Если источник пуст, удаление не происходит. Используется targetAlias для диамбигуации ссылок на столбцы. Доступно в Databricks Runtime 17.1 и более поздних версиях. Применяется к выборочному перезаписи данных с помощью Delta Lake. |
replaceUsing |
None | Разделенный запятыми список имен столбцов | Разделенный запятыми список имен столбцов, используемый для сопоставления строк между целевой таблицей и исходным запросом. Целевой объект и источник должны содержать все перечисленные столбцы. Строки в целевом объекте, соответствующие исходной строке при сравнении равенства, удаляются и заменяются.
NULL Значения обрабатываются как не равные и не будут соответствовать. Доступно в Databricks Runtime 16.3 и более поздних версиях. Применяется к выборочному перезаписи данных с помощью Delta Lake. |
replaceWhere |
None | Строка выражения предиката | Выражение предиката. Атомарно перезаписывает только записи, соответствующие предикату. Применяется к выборочному перезаписи данных с помощью Delta Lake. |
targetAlias |
None | Любая строка. | Псевдоним строки для целевой таблицы. Используйте для ссылки на столбцы или replaceOnreplaceWhere для диамбигуации, если условие ссылается на столбцы из целевой таблицы и исходного запроса. Применяется к выборочному перезаписи данных с помощью Delta Lake. |
txnAppId |
None | Любая строка. | Уникальная строка, определяющая приложение для идемпотентной записи в foreachBatch операциях. Используйте вместе с txnVersion тем, чтобы обеспечить точное запись в несколько таблиц Delta Lake. Применяется к использованию foreachBatch для записи идемпотентной таблицы. |
txnVersion |
None | Монотонно увеличивающееся целое число | Монотонно увеличивающееся число, используемое в качестве версии транзакции для идемпотентной записи в foreachBatch операциях. Используйте вместе с txnAppId тем, чтобы обеспечить точное запись в несколько таблиц Delta Lake. Применяется к использованию foreachBatch для записи идемпотентной таблицы. |
optimizeWrite |
None |
true, false |
Включение автоматической оптимизации записи для этой операции записи. Переопределяет конфигурацию spark.databricks.delta.optimizeWrite.enabled . Применяется к Несколько Delta Lake в Azure Databricks?. |
userMetadata |
None | Любая строка. | Определяемая пользователем строка, добавленная к метаданным фиксации для операции записи. Видимый в выходных DESCRIBE HISTORYданных . Применяется к таблицам Обогащения с пользовательскими метаданными. |
CSV
При написании CSV-файлов применяются следующие параметры.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
charToEscapeQuoteEscaping |
\0 (не включено) |
Один символ | Символ, используемый для escape-символа, если он отличается от символа кавычки. Применяется к csv-файлу (DataFrameWriter). |
compression |
none |
none (default), bzip2, gziplz4, snappy, deflate, zstd |
Кодек сжатия, используемый при записи. Применяется к csv-файлу (DataFrameWriter). |
dateFormat |
yyyy-MM-dd |
Строка формата даты | Формат строки для значений столбцов даты. Применяется к csv-файлу (DataFrameWriter). |
emptyValue |
Пустая строка | Любая строка. | Строка, написанная для пустых (непустых) значений. Применяется к csv-файлу (DataFrameWriter). |
encoding |
UTF-8 |
java.nio.charset.Charset Имя |
Кодировка символов для выходных файлов. Применяется к csv-файлу (DataFrameWriter). |
escape |
\ |
Один символ | Символ, используемый для escape-кавычек значений. Применяется к csv-файлу (DataFrameWriter). |
escapeQuotes |
true |
true, false |
Следует ли экранировать символы кавычки внутри значений поля с кавычками. Применяется к csv-файлу (DataFrameWriter). |
header |
false |
true, false |
Следует ли записывать имена столбцов в качестве первой строки выходных данных. Применяется к csv-файлу (DataFrameWriter). |
ignoreLeadingWhiteSpace |
false |
true, false |
Следует ли обрезать ведущие пробелы из значений при написании. Применяется к csv-файлу (DataFrameWriter). |
ignoreTrailingWhiteSpace |
false |
true, false |
Следует ли обрезать конечные пробелы из значений при записи. Применяется к csv-файлу (DataFrameWriter). |
lineSep |
\n |
Строка | Строка разделителя строк, используемая между записями. Применяется к csv-файлу (DataFrameWriter). |
locale |
en-US |
Идентификатор java.util.Locale |
Идентификатор java.util.Locale. Язык Java, который влияет на дату, метку времени и десятичный анализ по умолчанию в CSV. |
nullValue |
Пустая строка | Любая строка. | Строка, написанная для значений NULL. Применяется к csv-файлу (DataFrameWriter). |
quote |
" |
Один символ | Символ, используемый для значений полей кавычки, содержащих разделитель. Применяется к csv-файлу (DataFrameWriter). |
quoteAll |
false |
true, false |
Следует ли заключать все значения полей в кавычки независимо от содержимого. Применяется к csv-файлу (DataFrameWriter). |
sep |
, |
Строка | Символ разделителя полей. Применяется к csv-файлу (DataFrameWriter). |
timestampFormat |
yyyy-MM-dd'T'HH:mm:ss[.SSS][XXX] |
Строка формата метки времени | Строка форматирования для значений столбцов метки времени. Применяется к csv-файлу (DataFrameWriter). |
timestampNTZFormat |
yyyy-MM-dd'T'HH:mm:ss[.SSS] |
Строка формата метки времени | Форматирование строки метки времени без значений столбцов часового пояса (TimestampNTZType). |
Excel
Следующие параметры применяются при записи Excel файлов.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
dataAddress |
None | Имя листа или строка ссылки на ячейку | Имя листа или начальная ячейка записи. Если опущено, записывается на лист с именем Sheet1 , начиная с ячейки A1. Принимает имя листа (SheetName) или ссылку на одну ячейку (SheetName!A1). Диапазоны ячеек не поддерживаются для операций записи. |
dateFormatInWrite |
yyyy-mm-dd |
Строка формата даты Excel | Excel строку формата ячейки, примененную к столбцам Date. Использует синтаксис формата Excel. |
headerRows |
0 |
0, 1 |
Следует ли записывать имена столбцов в качестве первой строки. |
timestampNTZFormat |
yyyy-mm-dd hh:mm:ss |
Строка формата метки времени Excel | строка формата ячейки Excel применена к столбцам TimestampNTZ и Timestamp. Использует синтаксис формата Excel. |
version |
xlsx |
xlsx, xls |
Версия формата файла Excel для записи. |
JSON
При написании JSON-файлов применяются следующие параметры.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
compression |
none |
none, bzip2, gziplz4, snappy, deflate, zstd |
Кодек сжатия, используемый при записи. Применяется к json (DataFrameWriter). |
dateFormat |
yyyy-MM-dd |
Строка формата даты | Формат строки для значений столбцов даты. Применяется к json (DataFrameWriter). |
encoding |
UTF-8 |
java.nio.charset.Charset Имя |
Кодировка символов для выходных файлов. Применяется к json (DataFrameWriter). |
ignoreNullFields |
значение spark.sql.jsonGenerator.ignoreNullFields |
true, false |
Следует ли пропускать поля со значениями NULL из выходных данных JSON. Применяется к json (DataFrameWriter). |
lineSep |
\n |
Строка | Строка разделителя строк, используемая между записями. Применяется к json (DataFrameWriter). |
locale |
en-US |
Идентификатор java.util.Locale |
Java идентификатор языкового стандарта, влияющий на дату, метку времени и десятичный анализ в формате JSON. |
pretty |
false |
true, false |
Следует ли включать довольно (отступные, многострочный) выходные данные JSON. |
sortKeys |
false |
true, false |
Следует ли отсортировать ключи объектов JSON в алфавитном порядке в выходных данных. Полезно для производства детерминированных выходных данных. |
timestampFormat |
yyyy-MM-dd'T'HH:mm:ss[.SSS][XXX] |
Строка формата метки времени | Строка форматирования для значений столбцов метки времени. Применяется к json (DataFrameWriter). |
timestampNTZFormat |
yyyy-MM-dd'T'HH:mm:ss[.SSS] |
Строка формата метки времени | Форматирование строки метки времени без значений столбцов часового пояса (TimestampNTZType). |
writeNonAsciiCharacterAsCodePoint |
false |
true, false |
Следует ли кодировать символы, отличные от ASCII, как \uXXXX escape-последовательности Юникода вместо символов UTF-8 в выходных данных. |
ОРК
При написании файлов ORC применяются следующие параметры.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
compression |
zstd |
none, uncompressed, snappyzliblzozstdlz4brotli |
Кодек сжатия, используемый при записи. Применяется к orc (DataFrameWriter). |
Паркетным
При написании файлов Parquet применяются следующие параметры.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
compression |
snappy |
none, uncompressedsnappygziplzobrotlilz4lz4_rawzstd |
Кодек сжатия, используемый при записи. Применяется к parquet (DataFrameWriter). |
spark.sql.parquet.outputTimestampType |
INT96 |
INT96, , TIMESTAMP_MICROSTIMESTAMP_MILLIS |
Физический тип, используемый для кодирования столбцов метки времени. Используйте INT96 для совместимости с устаревшими средствами чтения Parquet, которые не поддерживают стандартные типы меток времени. |
Текст
При написании текстовых файлов применяются следующие параметры.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
compression |
none |
none, bzip2, gziplz4, snappy, deflate, zstd |
Кодек сжатия, используемый при записи. Применяется к тексту (DataFrameWriter). |
encoding |
UTF-8 |
java.nio.charset.Charset Имя |
Кодировка символов для выходных файлов. |
lineSep |
\n |
Строка | Строка разделителя строк, используемая между записями. Применяется к тексту (DataFrameWriter). |
XML
При написании XML-файлов применяются следующие параметры.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
arrayElementName |
item |
Любая строка. | Имя элемента для элементов массива, не имеющих явного имени. Применяется к xml (DataFrameWriter). |
attributePrefix |
_ |
Любая строка. | Префикс, подготовленный к именам полей, соответствующим атрибутам XML. Применяется к xml (DataFrameWriter). |
compression |
none |
none, bzip2, gziplz4, snappy, deflate, zstd |
Кодек сжатия, используемый при записи. Применяется к xml (DataFrameWriter). |
dateFormat |
yyyy-MM-dd |
Строка формата даты | Формат строки для значений столбцов даты. Применяется к xml (DataFrameWriter). |
declaration |
version="1.0" encoding="UTF-8" standalone="yes" |
Строка объявления XML или пустая строка для подавления | Строка объявления XML, написанная в верхней части каждого выходного файла. Задайте пустую строку, чтобы отключить объявление. Применяется к xml (DataFrameWriter). |
encoding |
UTF-8 |
java.nio.charset.Charset Имя |
Кодировка символов для выходных файлов. Применяется к xml (DataFrameWriter). |
indent |
4 пробела | Любая строка. | Строка, используемая для отступа дочерних элементов в выходных данных. Задайте пустую строку, чтобы отключить отступ и записать каждую строку в одной строке. |
locale |
en-US |
Идентификатор java.util.Locale |
Java идентификатор языкового стандарта, влияющий на дату, метку времени и десятичное форматирование в XML. |
nullValue |
null |
Любая строка. | Строка, написанная для значений NULL. Если задано значение null, атрибуты и дочерние элементы для полей NULL опущены. Применяется к xml (DataFrameWriter). |
rootTag |
ROWS |
Любая строка. | Тег корневого элемента, который упаковывает все элементы строки в выходные данные. Применяется к xml (DataFrameWriter). |
rowTag |
ROW |
Любая строка. | Тег элемента, представляющий строку в выходных данных. Применяется к xml (DataFrameWriter). |
singleVariantColumn |
None | Строка имени столбца | Имя одного столбца Variant для записи в XML-файлы. Применяется к xml (DataFrameWriter). |
timestampFormat |
yyyy-MM-dd'T'HH:mm:ss[.SSS][XXX] |
Строка формата метки времени | Строка форматирования для значений столбцов метки времени. Применяется к xml (DataFrameWriter). |
timestampNTZFormat |
yyyy-MM-dd'T'HH:mm:ss[.SSS] |
Строка формата метки времени | Форматирование строки метки времени без значений столбцов часового пояса. Применяется к xml (DataFrameWriter). |
validateName |
true |
true, false |
Следует ли вызывать исключение, если имя столбца не является допустимым идентификатором XML-элемента. Применяется к xml (DataFrameWriter). |
valueTag |
_VALUE |
Любая строка. | Имя поля, используемое для символьных данных в XML-элементах, которые также имеют атрибуты или дочерние элементы. Применяется к xml (DataFrameWriter). |
Параметры DataStreamWriter
Используйте эти параметры для DataStreamWriter.option() настройки потоковой записи.
Example
В следующем примере устанавливается расположение контрольной точки для потока:
Python
(df.writeStream
.format("delta")
.option("checkpointLocation", "/path/to/checkpoint")
.start("/path/to/table"))
Scala
df.writeStream
.format("delta")
.option("checkpointLocation", "/path/to/checkpoint")
.start("/path/to/table")
Общие
Следующие параметры применяются ко всем операциям потоковой записи.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
checkpointLocation |
Нет (обязательно) | Строка пути | Путь к каталогу контрольной точки для потокового запроса. Требуется для отказоустойчивости и точно один раз гарантии обработки. Каждый запрос потоковой передачи должен использовать уникальное расположение контрольной точки. Databricks рекомендует хранить контрольные точки в томе каталога Unity или пути к облачному хранилищу. Смотрите контрольные точки структурированной потоковой передачи. |
path |
None | Строка пути | Путь вывода для приемников потоковой передачи на основе файлов, таких как Parquet. Применяется только к форматам на основе файлов. |
Приемник консоли
Следующие параметры применяются при записи потоков в приемник консоли.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
numRows |
20 |
Положительные целые числа | Количество строк, отображаемых для каждого микропакета при записи в приемник консоли. |
truncate |
true |
true, false |
Следует ли усечь длинные строки при отображении строк. Установите для false отображения полных строковых значений. |
Delta Lake
Следующие параметры применяются при написании потока в таблицу Delta Lake с помощью format("delta"). Параметры, доступные только для перезаписи, например overwriteSchema, replaceWhereи partitionOverwriteMode не поддерживаются для потоковой записи.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
mergeSchema |
false |
true, false |
Следует ли развивать схему таблицы Delta Lake, когда потоковый кадр данных содержит новые столбцы. Применяется только к режиму добавления выходных данных. Применяется к схемам таблиц обновления с эволюцией схемы. |
userMetadata |
None | Любая строка. | Определяемая пользователем строка, добавленная к метаданным фиксации для операции записи. Видимый в выходных DESCRIBE HISTORYданных . Применяется к таблицам Обогащения с пользовательскими метаданными. |
Приемник файлов
Следующий параметр применяется при написании потока в форматы на основе файлов (Parquet, JSON, CSV, ORC, текст). Дополнительные сведения о параметрах формата см. в разделе "Параметры DataFrameWriter".
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
retention |
None | Строка времени, например 7 days или 24 hours |
Как долго сохранять файлы метаданных приемника, используемые для отказоустойчивости и сжатия. Если не задано, файлы метаданных сохраняются на неопределенный срок. |
Приемник Kafka
При написании в Kafka применяются следующие параметры.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
kafka.bootstrap.servers |
None | Разделенный запятыми список host:port строк |
Обязательно. Разделенный запятыми список адресов брокера host:port Kafka. |
topic |
None | Любая строка. | Целевой раздел Kafka для всех строк. Требуется, если кадр данных не содержит topic столбец. |
kafka.* |
None | Любое значение конфигурации производителя Kafka | Любая конфигурация производителя Kafka , префиксированная с kafka.префиксом . Например: kafka.compression.type. |
Приемник памяти
Следующие параметры применяются при записи потоков в приемник памяти.
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
queryName |
Нет (обязательно) | Любая строка. | Имя таблицы в памяти, в которую записывается запрос. Требуется для приемника памяти. Также можно настроить с помощью .queryName(). |
mode |
exactlyonce |
exactlyonce, atleastonce |
Гарантия доставки для приемника памяти.
exactlyonce использует микро-пакетный режим с точной семантикой.
atleastonce использует непрерывный режим с семантикой по крайней мере один раз. |
Параметры функции Spark
Некоторые встроенные функции Spark SQL принимают options карту, которая управляет анализом или сериализацией поведения. Передайте параметры в виде Python dict или Scala Map[String, String].
Example
В следующем примере анализируется столбец JSON при удалении неправильно сформированных записей:
Python
from pyspark.sql.functions import from_json
from pyspark.sql.types import StructType, StructField, StringType
schema = StructType([StructField("name", StringType())])
df = df.withColumn("parsed", from_json("json_col", schema, {"mode": "DROPMALFORMED"}))
Scala
import org.apache.spark.sql.functions.from_json
import org.apache.spark.sql.types._
val schema = StructType(Seq(StructField("name", StringType)))
val df = df.withColumn("parsed", from_json(col("json_col"), schema, Map("mode" -> "DROPMALFORMED")))
Avro
Функции Avro принимают те же параметры, что и соответствующие параметры кадра данных:
-
from_avroиспользуйтеschema_of_avroпараметры DataFrameReader Avro. -
to_avroиспользует параметры DataFrameWriter Avro.
Example
В следующем примере декодирует столбец Avro с включенной эволюцией схемы:
Python
from pyspark.sql.functions import from_avro
df = df.withColumn("decoded", from_avro("avro_col", json_schema, {"avroSchemaEvolutionMode": "restart"}))
Scala
import org.apache.spark.sql.avro.functions.from_avro
val df = df.withColumn("decoded", from_avro(col("avro_col"), jsonSchema, Map("avroSchemaEvolutionMode" -> "restart")))
Кроме того, варианты from_avro реестра схем и to_avro принимают следующие параметры:
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
schemaId |
None | Целое число идентификатора схемы | Идентификатор схемы из реестра схем Confluent, используемый при декодировании данных Avro, которые были закодированы с схемой, несовместимой с jsonFormatSchema. Применяется только к from_avro . |
confluent.schema.registry.* |
None | Любое значение свойства клиента Confluent SR | Свойства конфигурации клиента реестра схем Confluent. Передайте любое клиентское свойство Confluent SR с помощью этого префикса, например confluent.schema.registry.basic.auth.user.info для базовых учетных данных проверки подлинности. Требуется для вариантов и вариантов from_avroto_avroреестра схем. |
CSV
Функции CSV принимают те же параметры, что и соответствующие параметры кадра данных:
-
from_csvиспользуйтеschema_of_csvпараметры CSV DataFrameReader. -
to_csvиспользует параметры CSV-файла DataFrameWriter.
Example
Следующий пример считывает CSV-файл с пользовательским разделителем и NULL значением:
Python
from pyspark.sql.functions import from_csv
from pyspark.sql.types import StructType, StructField, IntegerType, StringType
schema = StructType([StructField("id", IntegerType()), StructField("name", StringType())])
df = df.withColumn("parsed", from_csv("csv_col", schema, {"sep": "|", "nullValue": "N/A"}))
Scala
import org.apache.spark.sql.functions.from_csv
import org.apache.spark.sql.types._
val schema = StructType(Seq(StructField("id", IntegerType), StructField("name", StringType)))
val df = df.withColumn("parsed", from_csv(col("csv_col"), schema, Map("sep" -> "|", "nullValue" -> "N/A")))
JSON
Функции JSON принимают те же параметры, что и соответствующие параметры кадра данных:
-
from_jsonиспользуйтеschema_of_jsonпараметры JSON DataFrameReader. -
to_jsonиспользует параметры JSON DataFrameWriter.
Example
В следующем примере выполняется запись JSON с полями NULL , которые игнорируются и включено довольно форматирование:
Python
from pyspark.sql.functions import to_json
df = df.withColumn("json_str", to_json("struct_col", {"pretty": "true", "ignoreNullFields": "true"}))
Scala
import org.apache.spark.sql.functions.to_json
val df = df.withColumn("json_str", to_json(col("struct_col"), Map("pretty" -> "true", "ignoreNullFields" -> "true")))
Protobuf
from_protobuf и to_protobuf не используйте файловый источник Данных. Данные Protobuf всегда считываются и записываются как двоичные столбцы с помощью этих функций. Параметры передаются в виде Map[String, String] регистра.
Example
В следующем примере декодирует столбец Protobuf с помощью режима PERMISSIVE:
Python
from pyspark.sql.functions import from_protobuf
df = df.withColumn("decoded", from_protobuf("proto_col", "MyMessage", "/path/to/descriptor.desc",
{"mode": "PERMISSIVE", "enums.as.ints": "true"}))
Scala
import org.apache.spark.sql.protobuf.functions.from_protobuf
val df = df.withColumn("decoded", from_protobuf(col("proto_col"), "MyMessage", "/path/to/descriptor.desc",
Map("mode" -> "PERMISSIVE", "enums.as.ints" -> "true")))
Функции Protobuf используют следующие параметры:
| Key | По умолчанию | Допустимые значения | Description |
|---|---|---|---|
mode |
FAILFAST |
FAILFAST, PERMISSIVE |
Обработка поврежденных записей.
FAILFAST выдает исключение.
PERMISSIVE Задает для неправильно сформированных полей значение NULL. Применимо к from_protobuf. |
recursive.fields.max.depth |
-1 (отключено) |
0 до 10 |
Максимальная глубина рекурсии для рекурсивных полей Protobuf. Установите для 0 отключения поддержки рекурсивных полей. Применимо к from_protobuf. |
convert.any.fields.to.json |
false |
true, false |
Следует ли преобразовывать поля Protobuf Any в строку JSON вместо строки STRUCT. Применимо к from_protobuf. |
emit.default.values |
false |
true, false |
Следует ли выдавать поля с нуля или значениями по умолчанию (семантика proto3). Если falseполя со значениями по умолчанию опущены из выходных данных. Применимо к from_protobuf. |
enums.as.ints |
false |
true, false |
Следует ли отображать поля перечисления в виде целых значений вместо строк. Применимо к from_protobuf. |
upcast.unsigned.ints |
false |
true, false |
Следует ли переадресовать uint32Long и uint64 предотвращать Decimal(20,0) переполнение целых чисел. Применимо к from_protobuf. |
unwrap.primitive.wrapper.types |
false |
true, false |
Следует ли распаковывать google.protobuf типы оболочки (например, Int32Value и StringValue) к соответствующим примитивным типам Spark. Применимо к from_protobuf. |
retain.empty.message.types |
false |
true, false |
Следует ли сохранять пустые типы сообщений Protobuf в выходной схеме, вставляя фиктивный столбец. Применимо к from_protobuf. |
schema.registry.subject |
None | Любая строка. | Имя субъекта реестра схем. Обязательный при использовании вариантов from_protobuf реестра схем и to_protobuf. |
schema.registry.address |
None |
host:port Строка |
Адрес реестра схем (узел и порт). Обязательный при использовании вариантов from_protobuf реестра схем и to_protobuf. |
schema.registry.protobuf.name |
None | Любая строка. | Указывает, какое сообщение Protobuf следует использовать, если тема реестра схем содержит несколько сообщений. Optional. |
schema.registry.schema.evolution.mode |
"restart" |
"restart", "none" |
Как изменения схемы обрабатываются при обнаружении более нового идентификатора схемы в входящей записи.
"restart" завершает запрос с UnknownFieldExceptionпомощью задания; настройте задания для перезапуска при сбое при выборе изменений.
"none" игнорирует изменения идентификатора схемы и анализирует более новые записи с исходной схемой. |
confluent.schema.registry.<option> |
— | Любое допустимое значение параметра клиента реестра схем Confluent | Передайте любой параметр реестра схем Confluent с помощью префикса "confluent.schema.registry". Например, задайте "confluent.schema.registry.basic.auth.credentials.source" значение "USER_INFO" и "confluent.schema.registry.basic.auth.user.info" настройте "<KEY>:<SECRET>" базовую проверку подлинности. |
XML
Функции XML принимают те же параметры, что и соответствующие параметры кадра данных:
-
from_xmlиспользуйтеschema_of_xmlпараметры XML DataFrameReader. -
to_xmlиспользует параметры XML DataFrameWriter.
Example
В следующем примере выполняется запись XML с пользовательскими тегами корневого и строкового кода:
Python
from pyspark.sql.functions import to_xml
df = df.withColumn("xml_str", to_xml("struct_col", {"rootTag": "records", "rowTag": "record"}))
Scala
import org.apache.spark.sql.functions.to_xml
val df = df.withColumn("xml_str", to_xml(col("struct_col"), Map("rootTag" -> "records", "rowTag" -> "record")))