Используйте mssql-python с Apache Arrow

Драйвер mssql-python предоставляет методы извлечения Apache Arrow для высокопроизводительного извлечения столбцевых данных из Microsoft SQL и База данных SQL Azure.

Apache Arrow — это платформа для кросс-языковой разработки колонковых данных в памяти. Драйвер преобразует наборы результатов ODBC напрямую в формат Arrow на C++, обходя создание объектов на Python для повышения производительности.

Интеграция со стрелкой позволяет:

  • Передача данных без копирования в Polars, pandas и DuckDB. «Zero-copy» означает, что данные остаются в одном буфере памяти, в который драйвер записывает данные, а использующие его библиотеки читают их напрямую, так что строки не копируются в промежуточные объекты Python.
  • Потоковая передача наборов результатов через RecordBatchReader без загрузки всего в память.
  • Формат столбцевых данных, идеально подходящий для аналитики и задач машинного обучения.
  • Снижение использования памяти по сравнению с построчным созданием объектов Python.

Методы курсора

Для использования методов выборки Arrow требуется пакет pyarrow. Установите это с pip install pyarrow. Если pyarrow он не установлен, вызов любого метода Arrow вызывает ImportError.

Драйвер mssql-python добавляет три метода к объекту курсора для доступа к данным Arrow. Все три метода преобразуют наборы результатов ODBC в формат Arrow на уровне C++ драйвера, что позволяет избежать создания промежуточных объектов Python.

  • arrow() возвращает весь набор результатов в виде одной таблицы в памяти. Самые простые в использовании.
  • arrow_batch() возвращает по одной порции строк за раз, что позволяет вам вручную управлять циклом.
  • arrow_reader() возвращает итератор, который автоматически выдает партии. Лучше всего подходит для потоковой передачи больших объёмов данных.

С использованием cursor.arrow(batch_size=8192)

Получите весь набор результатов в виде одного pyarrow.Table. Этот метод самый простой и хорошо работает, когда полный набор результатов помещается в память.

import mssql_python

conn = mssql_python.connect(connection_string)
cursor = conn.cursor()

cursor.execute("SELECT ProductID, Name, ListPrice FROM Production.Product")
table = cursor.arrow()

print(type(table))       # <class 'pyarrow.lib.Table'>
print(table.num_rows)    # Number of rows fetched
print(table.num_columns) # Number of columns
print(table.schema)      # Column names and Arrow types
print(table.to_pandas()) # Convert to pandas DataFrame

Замечание

Если в строке подключения используется Authentication=ActiveDirectoryDefault, драйвер использует DefaultAzureCredential, который поочередно проверяет несколько поставщиков учетных данных. Первое соединение может быть медленным, потому что SDK идёт по цепочке, пока не найдёт работающего провайдера. В продакшене, если вы знаете, какой тип учетных данных использует ваша среда, укажите его напрямую (например, ActiveDirectoryMSI для управляемой идентичности), чтобы избежать цепной ходьбы. Дополнительные сведения см. в разделе проверки подлинности Microsoft Entra.

С использованием cursor.arrow_batch(batch_size=8192)

Возьмите один pyarrow.RecordBatch , содержащий до batch_size строк. Используйте этот метод для кастомных циклов пакетной обработки, где нужен тонкий контроль над тем, сколько строк можно получить одновременно.

cursor.execute("SELECT * FROM Production.TransactionHistory")

while True:
    batch = cursor.arrow_batch(batch_size=10000)
    if batch.num_rows == 0:
        break
    # Process each batch
    print(f"Fetched {batch.num_rows} rows")

С использованием cursor.arrow_reader(batch_size=8192)

Возвращайте считыватель, который даёт RecordBatch объекты, пока набор результатов не исчерпается. Этот метод является наиболее эффективным по памяти вариантом для больших наборов результатов.

cursor.execute("SELECT * FROM Production.TransactionHistory")
reader = cursor.arrow_reader(batch_size=50000)

for batch in reader:
    # Process streaming batches without loading all data
    print(f"Batch: {batch.num_rows} rows")

Считыватель транслирует результаты по соединению, так что пока непрочитанный считыватель открыт, это соединение не может начать другое сообщение. Попытка выполнить это завершается ошибкой Connection is busy with results for another command.

Три вещи освобождают читателя: повторение до конца, закрытие родительского курсора или закрытие считывателя. Если вы прекратили чтение до того, как набор результатов будет полностью прочитан, и продолжаете использовать курсор, закройте объект чтения. При его закрытии также сбрасывается родительский курсор, чтобы на нём можно было выполнить другую инструкцию.

Используйте считыватель как менеджер контекста, чтобы он закрывался даже если исключение прерывает цикл:

cursor.execute("SELECT * FROM Production.TransactionHistory")

rows_seen = 0
with cursor.arrow_reader(batch_size=50000) as reader:
    for batch in reader:
        rows_seen += batch.num_rows
        if rows_seen >= 100000:
            break

# The reader is closed here, and the cursor is ready for the next statement.
cursor.execute("SELECT COUNT(*) FROM Production.TransactionHistory")

Вы также можете позвонить reader.close() напрямую. Вызывать его более одного раза безопасно, а свойство reader.closed показывает, закрыли ли вы его.

Распространенные шаблоны

Таблицы стрелок интегрируются напрямую с популярными библиотеками данных Python. Следующие примеры показывают, как передавать данные Arrow в pandas, Polars, DuckDB и в файлы различных форматов без копирования данных.

Загрузить результаты в pandas

cursor.execute("SELECT * FROM Production.Product")
table = cursor.arrow()

# Convert to pandas with zero-copy where possible
df = table.to_pandas()
print(df.head())

Загрузите результаты в Polars

import polars as pl

cursor.execute("SELECT * FROM Production.Product")
table = cursor.arrow()

df = pl.from_arrow(table)
print(df)

Результаты запросов с помощью DuckDB

DuckDB может отправлять запросы к таблицам Arrow непосредственно в SQL без копирования данных. Эта возможность полезна, когда нужен анализ в стиле SQL для наборов результатов, которые уже находятся в формате Arrow.

import duckdb

cursor.execute("SELECT * FROM Sales.SalesOrderHeader")
arrow_table = cursor.arrow()

# Query the Arrow table with DuckDB SQL
result = duckdb.sql("SELECT CustomerID, SUM(TotalDue) FROM arrow_table GROUP BY CustomerID")
print(result.fetchall())

Передавайте большие наборы результатов в формат Parquet

Для больших наборов результатов записывайте пакеты Arrow напрямую в файл Parquet в потоковом режиме, не загружая весь набор данных в память. ParquetWriter инкрементно записывает каждый пакет.

import pyarrow.parquet as pq

cursor.execute("SELECT * FROM Production.TransactionHistory")
reader = cursor.arrow_reader(batch_size=100000)

# Write streaming batches to a Parquet file
writer = None
for batch in reader:
    if writer is None:
        writer = pq.ParquetWriter("output.parquet", batch.schema)
    writer.write_batch(batch)

if writer:
    writer.close()

Экспорт в другие форматы

PyArrow предоставляет встроенные записчики для CSV и формата Arrow IPC (также известного как Feather V2). IPC-файлы Arrow точно сохраняют типы Arrow и быстро считываются обратно.

import pyarrow as pa
import pyarrow.csv as pcsv

cursor.execute("SELECT * FROM Production.Product")
table = cursor.arrow()

# Write to CSV
pcsv.write_csv(table, "products.csv")

# Write to an Arrow IPC file
with pa.ipc.new_file("products.arrow", table.schema) as writer:
    writer.write_table(table)

Загрузка данных Arrow в SQL Server

Метод cursor.bulkcopy_arrow() записывает данные Arrow в таблицу без преобразования её сначала в кортежи строк Python. Аргумент source принимает любое из следующих условий:

  • А pyarrow.Table.

  • А pyarrow.RecordBatch.

  • pyarrow.RecordBatchReader, включая средство чтения, возвращаемое cursor.arrow_reader().

  • Любой объект, который открывает интерфейс данных Arrow C через __arrow_c_stream__ или __arrow_c_array__.

import mssql_python
import pyarrow as pa

conn = mssql_python.connect(connection_string)

# bulkcopy_arrow() opens its own connection, so commit the table creation first.
conn.autocommit = True
cursor = conn.cursor()

cursor.execute("""
    CREATE TABLE ##SensorArchive (
        SensorID int NOT NULL,
        Reading float NULL,
        Location nvarchar(50) NULL
    )
""")

table = pa.table({
    "SensorID": pa.array([1, 2, 3], type=pa.int32()),
    "Reading": pa.array([20.5, None, 22.1], type=pa.float64()),
    "Location": pa.array(["Plant A", "Plant B", None], type=pa.string()),
})

result = cursor.bulkcopy_arrow("##SensorArchive", table)
print(f"Copied {result['rows_copied']} rows in {result['batch_count']} batches")

Значения NULL в Arrow записываются как значения SQL NULL.

Передать набор результатов в другую таблицу

Поскольку bulkcopy_arrow() принимает считыватель, вы можете перемещать большой набор результатов между таблицами, не материализуя его в памяти:

cursor.execute("""
    CREATE TABLE ##ProductArchive (
        ProductID int NOT NULL,
        Name nvarchar(50) NOT NULL,
        ListPrice money NOT NULL
    )
""")

cursor.execute("SELECT ProductID, Name, ListPrice FROM Production.Product")

with cursor.arrow_reader(batch_size=100000) as reader:
    result = cursor.bulkcopy_arrow("##ProductArchive", reader, batch_size=100000)

print(f"Copied {result['rows_copied']} rows")

Сопоставьте типы стрелок с столбцами назначения

Модуль записи Arrow требует, чтобы каждый тип столбца Arrow был совместим с соответствующим целевым типом столбца SQL. Она не конвертирует между семействами, поэтому несоответствие ValueError появляется до того, как записываются любые строки:

ValueError: Cannot map Arrow column 'ListPrice' (Float64) to SQL column 'ListPrice'
(Money): Usage Error: type combination is not supported by the Arrow row-major writer

Используйте отображения в отображениях типов данных в обратном направлении для выбора типа стрелки. money, decimal и numeric столбцы требуют decimal128, а не float64. Данные, считываемые обратно с помощью cursor.arrow(), уже имеют правильные типы, поэтому таблица, считываемая из SQL Server, загружается в соответствующую таблицу без преобразования.

Колонки карт по названию

Если порядок столбцов Arrow не совпадает с порядком столбцов в таблице назначения, передайте column_mappings, указав имена столбцов назначения в порядке столбцов Arrow:

from decimal import Decimal

table = pa.table({
    "Name": pa.array(["Widget"], type=pa.string()),
    "ProductID": pa.array([9001], type=pa.int32()),
    "ListPrice": pa.array([Decimal("12.34")], type=pa.decimal128(19, 4)),
})

cursor.bulkcopy_arrow(
    "##ProductArchive",
    table,
    column_mappings=["Name", "ProductID", "ListPrice"],
)

Метод принимает те же параметры, что и cursor.bulkcopy(), включая batch_size, timeout, keep_identity, table_lock и keep_nulls. Для получения дополнительной информации об этих вариантах см. раздел «Массовая копия».

Замечание

Передача источника Arrow в cursor.bulkcopy() вызывает TypeError и перенаправляет вас на cursor.bulkcopy_arrow().

Сопоставления типов данных

Методы извлечения Arrow сопоставляют типы Microsoft SQL с типами Arrow на уровне C++.

Microsoft SQL тип Тип стрелки
int, smallint, tinyint, bigint int32, int16, int8, int64
Плавающий, настоящий float64, float32
Десятичный, числовой decimal128
bit bool
Чар, Варчар, Нчар, Нварчар utf8
текст, ntext large_utf8
Бинарный, варбинарный binary, large_binary
date date32
time time64[us]
datetime, datetime2, smalldatetime timestamp[us]
datetimeoffset timestamp[us, tz=UTC]
uniqueidentifier utf8 (строка с заглавной буквы)
xml utf8

Замечание

Драйвер преобразует datetimeoffset тип в UTC, потому что столбцы стрелок требуют фиксированного часового пояса. Драйвер нормализует информацию о часовом поясе для каждой ячейки с Microsoft SQL в UTC во время конвертации.

Тип sql_variant не поддерживается методами выборки Arrow и вызывает исключение о неподдерживаемом типе данных. Используйте стандартный fetchone(), fetchmany(), или fetchall() для запросов, возвращающих sql_variant столбцы.

Вопросы производительности

Методы выборки стрелок наиболее быстры для аналитики и операций с массовыми данными, тогда как стандартные методы курсора лучше подходят для транзакционных паттернов с небольшими наборами результатов.

Когда использовать Arrow, а когда — стандартный fetch

Сценарий Рекомендуемый подход
Получить несколько строк для отображения fetchone() / fetchall()
Загрузите данные в Pandas или Polars cursor.arrow()
Обрабатывать большие наборы данных по частям cursor.arrow_reader()
Однострочные поиски или небольшие наборы результатов fetchone() / fetchval()
Аналитика или агрегационные конвейеры cursor.arrow() + Polars/DuckDB
Запишите результаты в Parquet или Arrow IPC cursor.arrow_reader() + PyArrow I/O

Управление памятью для больших наборов данных

Для наборов результатов, которые могут превышать доступную память, используйте arrow_reader() с разумным batch_size.

cursor.execute("SELECT * FROM Production.TransactionHistory")

# Process in batches of 100K rows
reader = cursor.arrow_reader(batch_size=100000)
total_rows = 0

for batch in reader:
    # Work with each batch individually
    total_rows += batch.num_rows
    # batch goes out of scope and memory is freed

print(f"Processed {total_rows} rows")

Настроить размер пакета

Параметр batch_size определяет, сколько строк будет получено в каждой партии. Оптимальный размер зависит от ширины строки и доступной памяти. Более широкие строки с большими столбцами, такими как nvarchar(max) или varbinary(max), лучше работают с меньшими размерами пакета, тогда как узкие строки — с большими.

  • По умолчанию (8192): Хороший баланс для большинства рабочих нагрузок.
  • Меньшие (1000-5000): используйте широкие столы с большими колонками.
  • Больше (50000-100000): используется для узких таблиц или когда пропускная способность важнее памяти.
# Narrow table with many rows - use larger batches
cursor.execute("SELECT ProductID, ListPrice FROM Production.Product")
table = cursor.arrow(batch_size=100000)

# Wide table with LOB columns - use smaller batches
cursor.execute("SELECT * FROM Production.Document")
table = cursor.arrow(batch_size=1000)