Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Polars — это высокопроизводительная библиотека DataFrame, написанная на Rust, которая предоставляет быструю и эффективную по памяти альтернативу pandas. Polars в сочетании с драйвером mssql-python позволяет вам:
- Загрузка результатов SQL-запроса напрямую в Polars DataFrames.
- Используйте Apache Arrow для передачи данных без копирования из Microsoft SQL.
- Эффективно записывайте Polars DataFrames обратно в Microsoft SQL.
- Создавайте высокопроизводительные конвейеры данных с ленивым анализом.
Примеры в этой статье используют запрос к AdventureWorks примерной базе данных. Если у вас его ещё нет, посмотрите примеры баз данных AdventureWorks.
Считывать данные в Polars DataFrames
Вы можете загрузить данные Microsoft SQL в Polars двумя способами: строка за строкой с помощью стандартных курсорных методов или передача без копирования через Apache Arrow. «Zero-copy» означает, что данные остаются в одном буфере памяти, из которого драйвер, Arrow и Polars читают напрямую, поэтому строки не дублируются в промежуточные объекты Python. Используйте подход Arrow для большинства рабочих нагрузок из-за такой эффективности.
Базовый запрос к DataFrame
Этот подход извлекает все строки с помощью стандартного курсора и вручную создаёт DataFrame Polars. Он принимает параметризованные запросы для безопасной замены значений. Он работает без PyArrow, но медленнее для больших наборов результатов, потому что каждое значение проходит через Python.
import polars as pl
import mssql_python
conn = mssql_python.connect(connection_string)
cursor = conn.cursor()
def query_to_polars(cursor, query: str, params: dict = None) -> pl.DataFrame:
"""Execute query and return results as Polars DataFrame."""
cursor.execute(query, params or {})
columns = [col[0] for col in cursor.description]
rows = cursor.fetchall()
data = {col: [row[i] for row in rows] for i, col in enumerate(columns)}
return pl.DataFrame(data)
# Usage: %(color)s is a parameterized placeholder. The driver safely substitutes
# the value from the dict, which prevents SQL injection.
df = query_to_polars(cursor, "SELECT TOP 5 Name, ListPrice FROM Production.Product WHERE Color = %(color)s", {"color": "Black"})
print(df)
Замечание
Если в строке подключения используется Authentication=ActiveDirectoryDefault, драйвер использует DefaultAzureCredential, который поочередно проверяет несколько поставщиков учетных данных. Первое соединение может быть медленным, потому что SDK идёт по цепочке, пока не найдёт работающего провайдера. В продакшене, если вы знаете, какой тип учетных данных использует ваша среда, укажите его напрямую (например, ActiveDirectoryMSI для управляемой идентичности), чтобы избежать цепной ходьбы. Дополнительные сведения см. в разделе проверки подлинности Microsoft Entra.
Используйте Arrow для передачи без копий (рекомендую)
Самый эффективный способ загрузки данных Microsoft SQL в Polars — это Apache Arrow. Метод arrow() драйвера mssql-python возвращает pyarrow.Table, который Polars может использовать без накладных расходов на копирование.
def query_to_polars_arrow(cursor, query: str, params: dict = None) -> pl.DataFrame:
"""Execute query and load results through Arrow for best performance."""
cursor.execute(query, params or {})
arrow_table = cursor.arrow()
return pl.from_arrow(arrow_table)
# Usage
df = query_to_polars_arrow(cursor, "SELECT ProductID, Name, ListPrice FROM Production.Product")
print(df)
Потоковая передача больших наборов данных с использованием пакетов Arrow
Для наборов данных, которые не помещаются в память, используйте arrow_reader() для обработки данных потоковыми пакетами. Каждая партия — это pyarrow.RecordBatch, который Polars может обрабатывать независимо, поэтому использование памяти остаётся пропорциональным batch_size, а не всему набору результатов.
def process_large_query(cursor, query: str, params: dict = None, batch_size: int = 50000) -> pl.DataFrame:
"""Process large query results as streaming Arrow batches."""
cursor.execute(query, params or {})
reader = cursor.arrow_reader(batch_size=batch_size)
results = []
for batch in reader:
chunk_df = pl.from_arrow(batch)
# Process each chunk
results.append(chunk_df)
return pl.concat(results) if results else pl.DataFrame()
# Usage
df = process_large_query(cursor, "SELECT * FROM Production.TransactionHistory")
Используйте LazyFrames для отложенного выполнения
Polars LazyFrames позволяет строить цепочку операций (фильтрация, группировка, сортировка) без запуска их сразу. Polars оптимизирует всю цепочку перед выполнением, что может быть быстрее, чем индивидуальное применение каждого шага.
def query_to_lazy(cursor, query: str, params: dict = None) -> pl.LazyFrame:
"""Execute query and return a Polars LazyFrame."""
cursor.execute(query, params or {})
arrow_table = cursor.arrow()
return pl.from_arrow(arrow_table).lazy()
# Build a query plan without executing immediately
lf = query_to_lazy(cursor, "SELECT SalesOrderID, CustomerID, TotalDue, OrderDate FROM Sales.SalesOrderHeader")
result = (
lf.filter(pl.col("TotalDue") > 100)
.group_by("CustomerID")
.agg([
pl.col("TotalDue").sum().alias("TotalSpent"),
pl.col("SalesOrderID").count().alias("OrderCount")
])
.sort("TotalSpent", descending=True)
.collect() # Execute the optimized plan
)
print(result)
Пишите Polars DataFrames в Microsoft SQL
Идентификаторы кавычок для предотвращения SQL-инъекции
Имена таблиц и столбцов не могут передаваться как параметры запроса в SQL. При создании SQL-выражений с динамическими идентификаторами заключайте каждое имя в квадратные скобки и экранируйте все встроенные символы ], чтобы предотвратить SQL-инъекции.
def quote_id(identifier: str) -> str:
"""Quote a Microsoft SQL identifier to prevent SQL injection.
Wraps the name in square brackets and escapes any embedded ] characters.
Raises ValueError if the identifier is empty or contains null bytes.
"""
if not identifier or "\x00" in identifier:
raise ValueError(f"Invalid identifier: {identifier!r}")
escaped = identifier.replace("]", "]]")
return f"[{escaped}]"
Вспомогательные функции в этом разделе используют quote_id() для всех имен таблиц и столбцов в сгенерированных SQL-запросах.
Вставить строки DataFrame
Построчный подход перебирает строки в DataFrame с помощью iter_rows(named=True) и выполняет один INSERT для каждой строки. Этот подход простой, но медленный для больших объёмов, потому что каждая строка требует обратной поездки на сервер.
def polars_to_sql(cursor, conn, df: pl.DataFrame, table: str) -> int:
"""Write Polars DataFrame to Microsoft SQL table."""
columns = df.columns
placeholders = ", ".join([f"%({col})s" for col in columns])
col_list = ", ".join([quote_id(col) for col in columns])
query = f"INSERT INTO {quote_id(table)} ({col_list}) VALUES ({placeholders})"
rows_inserted = 0
for row in df.iter_rows(named=True):
params = {k: (None if v is None else v) for k, v in row.items()}
cursor.execute(query, params)
rows_inserted += 1
conn.commit()
return rows_inserted
# Usage
cursor.execute("CREATE TABLE #PolarsInsert (Name NVARCHAR(100), Price DECIMAL(10,2), CategoryID INT)")
df = pl.DataFrame({
"Name": ["Product A", "Product B"],
"Price": [29.99, 49.99],
"CategoryID": [1, 2]
})
rows = polars_to_sql(cursor, conn, df, "#PolarsInsert")
print(f"Inserted {rows} rows")
Массовая вставка (рекомендуется для больших таблиц DataFrame)
Для больших DataFrames используйте метод драйвера bulkcopy(), чтобы пакетно отправлять строки по протоколу TDS (Tabular Data Stream) — собственному сетевому протоколу, используемому Microsoft SQL. Этот подход минимизирует круговые переходы и работает быстрее, чем вставки по строкам.
def polars_to_sql_bulk(conn, df: pl.DataFrame, table: str) -> int:
"""Bulk insert Polars DataFrame using BCP for best performance."""
rows = [tuple(None if v is None else v for v in row) for row in df.iter_rows()]
cursor = conn.cursor()
result = cursor.bulkcopy(table, rows)
conn.commit()
return result["rows_copied"]
# Usage
cursor.execute("CREATE TABLE ##PolarsBulk (Name NVARCHAR(50), Price FLOAT, CategoryID INT)")
conn.commit()
df = pl.DataFrame({
"Name": ["Product A", "Product B", "Product C"],
"Price": [29.99, 49.99, 19.99],
"CategoryID": [1, 2, 1]
})
rows = polars_to_sql_bulk(conn, df, "##PolarsBulk")
print(f"Bulk inserted {rows} rows")
Шаблоны анализа данных
Следующие примеры демонстрируют распространённые задачи анализа, сочетающие запросы Microsoft SQL с преобразованиями Polars.
Агрегированные запросы
Этот пример группирует продукты по подкатегориям и вычисляет статистику количества и цен в SQL, затем загружает сводку в Polars DataFrame:
def get_sales_summary(cursor) -> pl.DataFrame:
"""Get sales summary by subcategory."""
cursor.execute("""
SELECT
sc.Name AS SubcategoryName,
COUNT(*) AS ProductCount,
AVG(p.ListPrice) AS AvgPrice,
MIN(p.ListPrice) AS MinPrice,
MAX(p.ListPrice) AS MaxPrice
FROM Production.Product p
JOIN Production.ProductSubcategory sc ON p.ProductSubcategoryID = sc.ProductSubcategoryID
GROUP BY sc.Name
ORDER BY ProductCount DESC
""")
return pl.from_arrow(cursor.arrow())
df = get_sales_summary(cursor)
print(df)
Аналитика временных рядов
Загрузите временные ряды данных из Microsoft SQL и добавляйте вычисленные столбцы, такие как скользящие средние, используя выражения Polars.
def get_daily_sales(cursor, start_date: str, end_date: str) -> pl.DataFrame:
"""Get daily sales and compute rolling statistics."""
cursor.execute("""
SELECT
CAST(OrderDate AS DATE) AS Date,
COUNT(*) AS OrderCount,
SUM(TotalDue) AS Revenue
FROM Sales.SalesOrderHeader
WHERE OrderDate BETWEEN %(start)s AND %(end)s
GROUP BY CAST(OrderDate AS DATE)
ORDER BY Date
""", {"start": start_date, "end": end_date})
df = pl.from_arrow(cursor.arrow())
# Add rolling 7-day average
df = df.with_columns(
pl.col("Revenue").rolling_mean(window_size=7).alias("RollingAvg")
)
return df
sales_df = get_daily_sales(cursor, "2013-01-01", "2013-12-31")
print(sales_df)
Объединяйте данные SQL с локальными файлами
Вы можете обогатить данные Microsoft SQL, объединив их с локальными CSV-файлами в Polars. Загрузите каждый источник в DataFrame и объедините их в памяти.
# Load SQL data via Arrow
cursor.execute("SELECT c.CustomerID, p.FirstName, p.LastName FROM Sales.Customer c JOIN Person.Person p ON c.PersonID = p.BusinessEntityID")
customers = pl.from_arrow(cursor.arrow())
# Load local CSV
orders = pl.read_csv("orders_export.csv")
# Join in Polars
result = customers.join(orders, on="CustomerID", how="inner")
print(result)
Паттерны ETL
Постройте конвейеры для извлечения, преобразования и загрузки (ETL), комбинируя запросы Microsoft SQL с трансформациями Polars. Выражения Polars отвечают за этап преобразования, а bulkcopy() — за загрузку.
Извлечение, преобразование, загрузка
Этот пример извлекает данные об активных клиентах с помощью Arrow, применяет логику бизнес-сегментации с использованием выражений Polars и загружает результаты с помощью пакетной загрузки.
def etl_pipeline(source_cursor, dest_conn):
"""ETL pipeline using Polars transformations."""
# Extract: derive a per-customer summary from order history via Arrow
source_cursor.execute("""
SELECT
CustomerID,
COUNT(*) AS OrderCount,
SUM(TotalDue) AS TotalSpent
FROM Sales.SalesOrderHeader
WHERE OrderDate > DATEADD(YEAR, -1, (SELECT MAX(OrderDate) FROM Sales.SalesOrderHeader))
GROUP BY CustomerID
""")
df = pl.from_arrow(source_cursor.arrow())
df = df.with_columns(pl.col("TotalSpent").cast(pl.Float64))
# Transform with Polars expressions
df = df.with_columns([
pl.when(pl.col("TotalSpent") > 1000).then(pl.lit("Platinum"))
.when(pl.col("TotalSpent") > 500).then(pl.lit("Gold"))
.when(pl.col("TotalSpent") > 100).then(pl.lit("Silver"))
.otherwise(pl.lit("Bronze"))
.alias("CustomerSegment"),
(pl.col("TotalSpent") / pl.col("OrderCount").clip(lower_bound=1))
.alias("AvgOrderValue"),
(pl.col("TotalSpent") > 500).alias("IsHighValue")
])
# Load via bulk copy into the destination table
dest_cursor = dest_conn.cursor()
dest_cursor.execute("""
CREATE TABLE ##CustomerAnalytics (
CustomerID INT,
CustomerSegment NVARCHAR(20),
AvgOrderValue FLOAT,
IsHighValue BIT
)
""")
dest_conn.commit()
load_df = df.select(["CustomerID", "CustomerSegment", "AvgOrderValue", "IsHighValue"])
polars_to_sql_bulk(dest_conn, load_df, "##CustomerAnalytics")
return len(df)
Советы по производительности
Следующие советы помогут вам максимально эффективно использовать комбинацию mssql-python и Polars.
Пусть Microsoft SQL возьмёт на себя тяжёлую работу
Microsoft SQL работает быстрее для агрегирования, фильтрации и объединений, чем перемещение всех необработанных данных по проводу и локальную обработку в Python. Пусть Microsoft SQL выполняет основную работу, когда это возможно, перемещайте только нужные данные по сети и используйте Polars для анализа и трансформаций, которые удобнее в Python.
# Avoid: pulling all rows over the wire to aggregate locally in Polars
df = query_to_polars_arrow(cursor, "SELECT * FROM Production.Product WHERE Color IS NOT NULL") # transfers entire table
summary = df.group_by("Color").agg(pl.col("ListPrice").sum()) # aggregation that SQL can do faster
# Better: push the aggregation into SQL and transfer only the summary
df = query_to_polars_arrow(cursor, """
SELECT Color AS Category, SUM(ListPrice) AS TotalAmount
FROM Production.Product
WHERE Color IS NOT NULL
GROUP BY Color
""")
Используйте Arrow для всех операций чтения
Передача на основе стрелок предотвращает создание промежуточных объектов Python, что снижает использование памяти и повышает пропускную способность. Предпочитайте cursor.arrow(), а не ручное построчное преобразование для любого набора результатов, содержащего больше нескольких строк.
# Suboptimal: Row-by-row conversion
cursor.execute("SELECT * FROM Production.TransactionHistory")
columns = [col[0] for col in cursor.description]
rows = cursor.fetchall()
df = pl.DataFrame({col: [row[i] for row in rows] for i, col in enumerate(columns)})
# Better: Arrow-based transfer
cursor.execute("SELECT * FROM Production.TransactionHistory")
df = pl.from_arrow(cursor.arrow())