本文档介绍了 DataStore API 中的核心类。
DataStore
用于数据处理的主要 DataFrame 风格类。
from chdb.datastore import DataStore构造函数
DataStore(data=None, columns=None, index=None, dtype=None, copy=None)参数:
| 参数 | 类型 | 描述 |
|---|---|---|
data |
dict/list/DataFrame/DataStore | 输入数据 |
columns |
list | 列名 |
index |
Index | 行索引 |
dtype |
dict | 列数据类型 |
copy |
bool | 复制数据 |
示例:
# From dictionary
ds = DataStore({'a': [1, 2, 3], 'b': ['x', 'y', 'z']})
# From pandas DataFrame
import pandas as pd
ds = DataStore(pd.DataFrame({'a': [1, 2, 3]}))
# Empty DataStore
ds = DataStore()属性
| 属性 | 类型 | 描述 |
|---|---|---|
columns |
Index | 列名 |
dtypes |
Series | 各列的数据类型 |
shape |
tuple | (行数, 列数) |
size |
int | 元素总数 |
ndim |
int | 维度数 (2) |
empty |
bool | DataFrame 是否为空 |
values |
ndarray | 底层数据的 NumPy 数组表示 |
index |
Index | 行索引 |
T |
DataStore | 转置 |
axes |
list | 轴的列表 |
工厂方法
| 方法 | 说明 |
|---|---|
uri(uri) |
通用 URI 工厂方法 |
from_file(path, ...) |
从文件创建 |
from_df(df) |
从 pandas DataFrame 创建 |
from_s3(url, ...) |
从 S3 创建 |
from_gcs(url, ...) |
从 Google Cloud Storage 创建 |
from_azure(url, ...) |
从 Azure Blob 创建 |
from_mysql(...) |
从 MySQL 创建 |
from_postgresql(...) |
从 PostgreSQL 创建 |
from_clickhouse(...) |
从 ClickHouse 创建 |
from_mongodb(...) |
从 MongoDB 创建 |
from_sqlite(...) |
从 SQLite 创建 |
from_iceberg(path) |
从 Iceberg 表创建 |
from_delta(path) |
从 Delta Lake 创建 |
from_numbers(n) |
使用连续数字创建 |
from_random(rows, cols) |
使用随机数据创建 |
run_sql(query) |
从 SQL 查询创建 |
详见工厂方法。
查询方法
| 方法 | 返回值 | 描述 |
|---|---|---|
select(*cols) |
DataStore | 选择列 |
filter(condition) |
DataStore | 过滤行 |
where(condition) |
DataStore | filter 的别名 |
sort(*cols, ascending=True) |
DataStore | 对行进行排序 |
orderby(*cols) |
DataStore | sort 的别名 |
limit(n) |
DataStore | 限制行数 |
offset(n) |
DataStore | 跳过若干行 |
distinct(subset=None) |
DataStore | 去除重复项 |
groupby(*cols) |
LazyGroupBy | 对行进行分组 |
having(condition) |
DataStore | 过滤分组结果 |
join(right, ...) |
DataStore | 连接 DataStore |
union(other, all=False) |
DataStore | 合并 DataStore |
when(cond, val) |
CaseWhen | CASE WHEN |
详见查询构建。
Pandas 兼容方法
完整的 209 个方法列表,请参阅 Pandas 兼容性。
索引:
head(), tail(), sample(), loc, iloc, at, iat, query(), isin(), where(), mask(), get(), xs(), pop()
聚合:
sum(), mean(), std(), var(), min(), max(), median(), count(), nunique(), quantile(), describe(), corr(), cov(), skew(), kurt()
操作:
drop(), drop_duplicates(), dropna(), fillna(), replace(), rename(), assign(), astype(), copy()
排序:
sort_values(), sort_index(), nlargest(), nsmallest(), rank()
重塑:
pivot(), pivot_table(), melt(), stack(), unstack(), transpose(), explode(), squeeze()
组合:
merge(), join(), concat(), append(), combine(), update(), compare()
应用/转换:
apply(), applymap(), map(), agg(), transform(), pipe(), groupby()
时间序列:
rolling(), expanding(), ewm(), shift(), diff(), pct_change(), resample()
I/O 方法
| 方法 | 描述 |
|---|---|
to_csv(path, ...) |
导出为 CSV |
to_parquet(path, ...) |
导出为 Parquet |
to_json(path, ...) |
导出为 JSON |
to_excel(path, ...) |
导出为 Excel |
to_df() |
转换为 pandas DataFrame |
to_pandas() |
to_df 的别名 |
to_arrow() |
转换为 Arrow Table |
to_dict(orient) |
转换为字典 |
to_records() |
转换为记录 |
to_numpy() |
转换为 NumPy 数组 |
to_sql() |
生成 SQL 字符串 |
to_string() |
字符串表示形式 |
to_markdown() |
Markdown 表格 |
to_html() |
HTML 表格 |
详见 I/O 操作。
调试方法
| 方法 | 描述 |
|---|---|
explain(verbose=False) |
查看执行计划 |
clear_cache() |
清除缓存结果 |
详情请参见调试。
魔术方法
| 方法 | 描述 | |
|---|---|---|
__getitem__(key) |
ds['col'], ds[['a', 'b']], ds[condition] |
|
__setitem__(key, value) |
ds['col'] = value |
|
__delitem__(key) |
del ds['col'] |
|
__len__() |
len(ds) |
|
__iter__() |
for col in ds |
|
__contains__(key) |
'col' in ds |
|
__repr__() |
repr(ds) |
|
__str__() |
str(ds) |
|
__eq__(other) |
ds == other |
|
__ne__(other) |
ds != other |
|
__lt__(other) |
ds < other |
|
__le__(other) |
ds <= other |
|
__gt__(other) |
ds > other |
|
__ge__(other) |
ds >= other |
|
__add__(other) |
ds + other |
|
__sub__(other) |
ds - other |
|
__mul__(other) |
ds * other |
|
__truediv__(other) |
ds / other |
|
__floordiv__(other) |
ds // other |
|
__mod__(other) |
ds % other |
|
__pow__(other) |
ds ** other |
|
__and__(other) |
ds & other |
|
__or__(other) |
`ds | other` |
__invert__() |
~ds |
|
__neg__() |
-ds |
|
__pos__() |
+ds |
|
__abs__() |
abs(ds) |
ColumnExpr
表示用于惰性求值的列表达式。在访问列时会返回该表达式。
# ColumnExpr is returned automatically
col = ds['name'] # Returns ColumnExpr属性
| 属性 | 类型 | 描述 |
|---|---|---|
name |
str | 列名 |
dtype |
dtype | 数据类型 |
访问器
| 访问器 | 描述 | 方法 |
|---|---|---|
.str |
字符串操作 | 56 个方法 |
.dt |
日期时间操作 | 42+ 个方法 |
.arr |
数组操作 | 37 个方法 |
.json |
JSON 解析 | 13 个方法 |
.url |
URL 解析 | 15 个方法 |
.ip |
IP 地址操作 | 9 个方法 |
.geo |
地理空间/距离操作 | 14 个方法 |
完整文档请参见 访问器。
算术操作
ds['total'] = ds['price'] * ds['quantity']
ds['profit'] = ds['revenue'] - ds['cost']
ds['ratio'] = ds['a'] / ds['b']
ds['squared'] = ds['value'] ** 2
ds['remainder'] = ds['value'] % 10比较运算
ds[ds['age'] > 25] # Greater than
ds[ds['age'] >= 25] # Greater or equal
ds[ds['age'] < 25] # Less than
ds[ds['age'] <= 25] # Less or equal
ds[ds['name'] == 'Alice'] # Equal
ds[ds['name'] != 'Bob'] # Not equal逻辑操作
ds[(ds['age'] > 25) & (ds['city'] == 'NYC')] # AND
ds[(ds['age'] > 25) | (ds['city'] == 'NYC')] # OR
ds[~(ds['status'] == 'inactive')] # NOT方法
| 方法 | 描述 |
|---|---|
as_(alias) |
设置别名 |
cast(dtype) |
转换为指定类型 |
astype(dtype) |
cast 的别名 |
isnull() |
是否为 NULL |
notnull() |
是否不为 NULL |
isna() |
isnull 的别名 |
notna() |
notnull 的别名 |
isin(values) |
是否在值列表中 |
between(low, high) |
是否介于两个值之间 |
fillna(value) |
填充 NULL 值 |
replace(to_replace, value) |
替换值 |
clip(lower, upper) |
裁剪值 |
abs() |
取绝对值 |
round(decimals) |
对值进行四舍五入 |
floor() |
向下取整 |
ceil() |
向上取整 |
apply(func) |
应用函数 |
map(mapper) |
映射值 |
聚合方法
| 方法 | 描述 |
|---|---|
sum() |
求和 |
mean() |
均值 |
avg() |
mean 的别名 |
min() |
最小值 |
max() |
最大值 |
count() |
统计非 NULL 值 |
nunique() |
去重计数 |
std() |
标准差 |
var() |
方差 |
median() |
中位数 |
quantile(q) |
分位数 |
first() |
第一个值 |
last() |
最后一个值 |
any() |
任一值为 true |
all() |
全部为 true |
LazyGroupBy
表示一个可用于聚合操作的分组 DataStore。
# LazyGroupBy is returned automatically
grouped = ds.groupby('category') # Returns LazyGroupBy方法
| 方法 | 返回值 | 描述 |
|---|---|---|
agg(spec) |
DataStore | 聚合 |
aggregate(spec) |
DataStore | agg 的别名 |
sum() |
DataStore | 各组求和 |
mean() |
DataStore | 各组均值 |
count() |
DataStore | 各组计数 |
min() |
DataStore | 各组最小值 |
max() |
DataStore | 各组最大值 |
std() |
DataStore | 各组标准差 |
var() |
DataStore | 各组方差 |
median() |
DataStore | 各组中位数 |
nunique() |
DataStore | 各组唯一值计数 |
first() |
DataStore | 各组的第一个值 |
last() |
DataStore | 各组的最后一个值 |
nth(n) |
DataStore | 各组的第 n 个值 |
head(n) |
DataStore | 各组前 n 个值 |
tail(n) |
DataStore | 各组后 n 个值 |
apply(func) |
DataStore | 对各组应用函数 |
transform(func) |
DataStore | 对各组进行转换 |
filter(func) |
DataStore | 过滤各组 |
列选择
# Select column after groupby
grouped['amount'].sum() # Returns DataStore
grouped[['a', 'b']].sum() # Returns DataStore聚合规范
# Single aggregation
grouped.agg({'amount': 'sum'})
# Multiple aggregations per column
grouped.agg({'amount': ['sum', 'mean', 'count']})
# Named aggregations
grouped.agg(
total=('amount', 'sum'),
average=('amount', 'mean'),
count=('id', 'count')
)LazySeries
表示惰性 Series (单列) 。
属性
| 属性 | 类型 | 描述 |
|---|---|---|
name |
str | 序列名称 |
dtype |
dtype | 数据类型 |
方法
继承自 ColumnExpr 的大多数方法。主要方法:
| 方法 | 描述 |
|---|---|
value_counts() |
值出现频率 |
unique() |
去重后的值 |
nunique() |
唯一值数量 |
mode() |
众数 |
to_list() |
转换为列表 |
to_numpy() |
转换为数组 |
to_frame() |
转换为 DataStore |
F (函数)
用于 ClickHouse 函数的命名空间。
from chdb.datastore import F, Field
# Aggregations
F.sum(Field('amount'))
F.avg(Field('price'))
F.count(Field('id'))
F.quantile(Field('value'), 0.95)
# Conditional
F.sum_if(Field('amount'), Field('status') == 'completed')
F.count_if(Field('active'))
# Window
F.row_number().over(order_by='date')
F.lag('price', 1).over(partition_by='product', order_by='date')详见聚合。
字段
按名称引用列。
from chdb.datastore import Field
# Create field reference
amount = Field('amount')
price = Field('price')
# Use in expressions
F.sum(Field('amount'))
F.avg(Field('price'))CaseWhen
CASE WHEN 表达式构建器。
# Create case-when expression
result = (ds
.when(ds['score'] >= 90, 'A')
.when(ds['score'] >= 80, 'B')
.when(ds['score'] >= 70, 'C')
.otherwise('F')
)
# Assign to column
ds['grade'] = resultWindow
窗口函数的窗口规范。
from chdb.datastore import F
# Create window
window = F.window(
partition_by='category',
order_by='date',
rows_between=(-7, 0)
)
# Use with aggregation
ds['rolling_avg'] = F.avg('price').over(window)