Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

DataStore 类参考

本文档介绍了 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'] = result

Window

窗口函数的窗口规范。

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)
Navigation