--- name: zarr-python description: 用于云存储的分块N维数组(Zarr-Python 3)。压缩数组、并行I/O、通过fsspec的S3/GCS、兼容NumPy/Dask/Xarray,用于大规模科学计算管道。 allowed-tools: Read Write Edit Bash license: MIT license compatibility: Requires Python 3.12+ and zarr 3.x. Cloud I/O needs zarr[remote] plus pinned s3fs or gcsfs. Legacy Zarr v2 workflows need exact 2.x pins on older Python. metadata: {"version": "1.1", "skill-author": "K-Dense Inc."} --- # Zarr Python ## 概述 Zarr是一个Python库,用于存储带有分块和压缩的大型N维数组。应用此技能可实现高效并行I/O、云原生工作流,以及与NumPy、Dask和Xarray的无缝集成。 **当前上游版本:** zarr **3.2.1**(发布于2026-05-05)。文档:[zarr.readthedocs.io](https://zarr.readthedocs.io/en/stable/)。新数组默认使用**Zarr format 3**;如需兼容旧版可设置`zarr_format=2`。Zarr 3.2 增加了矩形分块(rectilinear chunks),并持续完善v3编解码器流水线。此技能是由 K-Dense Inc. 维护的**社区指南**,并非 zarr-developers 官方软件包。 ## 快速开始 ### 安装 ```bash uv pip install "zarr==3.2.1" ``` 当前稳定版 Zarr-Python 需要 **Python 3.12+** 和 NumPy 2.0+。对于远程存储(S3、GCS、HTTP),请在项目锁定文件中固定可选的扩展/后端版本: ```bash uv pip install "zarr[remote]==3.2.1" "s3fs==2026.4.0" "gcsfs==2026.5.0" ``` 仅当项目已有提交的锁定文件和兼容性测试时,才使用如 `zarr>=3,<4` 这样的版本范围。对于 Zarr-Python 2 / Python 3.10–3.11 的工作流,请从 support-v2 发行说明中选择一个精确的 `zarr==2.x.y` 补丁版本,并提交生成的锁定文件。 ### 基本数组创建 ```python import zarr import numpy as np # 创建带有分块和压缩的2D数组 z = zarr.create_array( store="data/my_array.zarr", shape=(10000, 10000), chunks=(1000, 1000), dtype="f4" ) # 使用NumPy风格索引写入数据 z[:, :] = np.random.random((10000, 10000)) # 读取数据 data = z[0:100, 0:100] # 返回NumPy数组 ``` ## 核心操作 ### 创建数组 Zarr提供多个方便的数组创建函数: ```python # 创建空数组 z = zarr.zeros(shape=(10000, 10000), chunks=(1000, 1000), dtype='f4', store='data.zarr') # 创建填充数组 z = zarr.ones((5000, 5000), chunks=(500, 500)) z = zarr.full((1000, 1000), fill_value=42, chunks=(100, 100)) # 从现有数据创建 data = np.arange(10000).reshape(100, 100) z = zarr.array(data, chunks=(10, 10), store='data.zarr') # 创建与另一个数组相似的数组 z2 = zarr.zeros_like(z) # 匹配z的形状、分块和数据类型 ``` ### 打开现有数组 ```python # 打开数组(默认为读写模式) z = zarr.open_array('data.zarr', mode='r+') # 只读模式 z = zarr.open_array('data.zarr', mode='r') # open()函数自动检测数组与组 z = zarr.open('data.zarr') # 返回Array或Group ``` ### 读写数据 Zarr数组支持NumPy风格索引: ```python # 写入整个数组 z[:] = 42 # 写入切片 z[0, :] = np.arange(100) z[10:20, 50:60] = np.random.random((10, 10)) # 读取数据(返回NumPy数组) data = z[0:100, 0:100] row = z[5, :] # 高级索引 z.vindex[[0, 5, 10], [2, 8, 15]] # 坐标索引 z.oindex[0:10, [5, 10, 15]] # 正交索引 z.blocks[0, 0] # 块/分块索引 ``` ### 调整大小和追加 ```python # 调整数组大小(v3:以元组形式传入形状) z.resize((15000, 15000)) # 沿轴追加数据 z.append(np.random.random((1000, 10000)), axis=0) # 添加行 ``` ## 分块策略 分块对性能至关重要。根据访问模式选择分块大小和形状。 ### 分块大小指南 - **最小分块大小**:推荐1 MB以获得最佳性能 - **平衡**:较大的分块 = 较少的元数据操作;较小的分块 = 更好的并行访问 - **内存考虑**:压缩期间整个分块必须适合内存 ```python # 配置分块大小(目标为每个分块~1MB) # 对于float32数据:1MB = 262,144个元素 = 512×512数组 z = zarr.zeros( shape=(10000, 10000), chunks=(512, 512), # ~1MB分块 dtype='f4' ) ``` ### 使分块与访问模式对齐 **关键**:分块形状根据数据访问方式显著影响性能。 ```python # 如果频繁访问行(第一维度) z = zarr.zeros((10000, 10000), chunks=(10, 10000)) # 分块跨越列 # 如果频繁访问列(第二维度) z = zarr.zeros((10000, 10000), chunks=(10000, 10)) # 分块跨越行 # 对于混合访问模式(平衡方法) z = zarr.zeros((10000, 10000), chunks=(1000, 1000)) # 方形分块 ``` **性能示例**:对于(200, 200, 200)数组,沿第一维度读取: - 使用分块(1, 200, 200):~107ms - 使用分块(200, 200, 1):~1.65ms(快65倍!) ### 矩形分块与分片 Zarr 3.2 支持**矩形分块(rectilinear chunks)**,用于处理不均匀的网格。当某个维度存在大小不一的分片时,可以传入嵌套的分块长度: ```python z = zarr.create_array( store="rectilinear.zarr", shape=(60, 100), chunks=([10, 20, 30], [50, 50]), dtype="f4", ) ``` 当数组有数百万个小分块时,使用**分片(sharding)**将分块分组为更大的存储对象: ```python # 创建带分片的数组 z = zarr.create_array( store='data.zarr', shape=(100000, 100000), chunks=(100, 100), # 小分块用于访问 shards=(1000, 1000), # 每个分片组100个分块 dtype='f4' ) ``` **好处**: - 减少来自数百万小文件的文件系统开销 - 提高云存储性能(减少对象请求) - 防止文件系统块大小浪费 **重要**:写入前整个分片必须适合内存。 ## 压缩 Zarr按分块应用压缩以减少存储同时保持快速访问。 ### 配置压缩 ```python from zarr.codecs import BloscCodec, BloscShuffle, GzipCodec # 默认:Blosc与Zstandard z = zarr.zeros((1000, 1000), chunks=(100, 100)) # 使用默认压缩 # 配置Blosc压缩 z = zarr.create_array( store='data.zarr', shape=(1000, 1000), chunks=(100, 100), dtype='f4', compressors=BloscCodec(cname='zstd', clevel=5, shuffle=BloscShuffle.bitshuffle) ) # 可用的Blosc压缩器:'blosclz', 'lz4', 'lz4hc', 'snappy', 'zlib', 'zstd' # 使用Gzip压缩 z = zarr.create_array( store='data.zarr', shape=(1000, 1000), chunks=(100, 100), dtype='f4', compressors=GzipCodec(level=6) ) # 禁用压缩 z = zarr.create_array( store='data.zarr', shape=(1000, 1000), chunks=(100, 100), dtype='f4', compressors=None ) ``` ### 压缩性能提示 - **Blosc**(默认):快速压缩/解压缩,适合交互式工作负载 - **Zstandard**:更好的压缩率,比LZ4稍慢 - **Gzip**:最大压缩,性能较慢 - **LZ4**:最快压缩,较低压缩率 - **Shuffle**:对数值数据启用shuffle过滤器以获得更好的压缩 ```python # 对数值科学数据最佳 compressors=BloscCodec(cname='zstd', clevel=5, shuffle=BloscShuffle.bitshuffle) # 对速度最佳 compressors=BloscCodec(cname='lz4', clevel=1) # 对压缩率最佳 compressors=GzipCodec(level=9) ``` ## 存储后端 Zarr支持多种存储后端通过灵活的存储接口。 ### 本地文件系统(默认) ```python from zarr.storage import LocalStore # 显式创建存储 store = LocalStore('data/my_array.zarr') z = zarr.open_array(store=store, mode='w', shape=(1000, 1000), chunks=(100, 100)) # 或使用字符串路径(自动创建LocalStore) z = zarr.open_array('data/my_array.zarr', mode='w', shape=(1000, 1000), chunks=(100, 100)) ``` ### 内存存储 ```python from zarr.storage import MemoryStore # 创建内存存储 store = MemoryStore() z = zarr.open_array(store=store, mode='w', shape=(1000, 1000), chunks=(100, 100)) # 数据仅存在于内存中,不持久化 ``` ### ZIP文件存储 ```python from zarr.storage import ZipStore # 写入ZIP文件 store = ZipStore('data.zip', mode='w') z = zarr.open_array(store=store, mode='w', shape=(1000, 1000), chunks=(100, 100)) z[:] = np.random.random((1000, 1000)) store.close() # 重要:必须关闭ZipStore # 从ZIP文件读取 store = ZipStore('data.zip', mode='r') z = zarr.open_array(store=store) data = z[:] store.close() ``` ### 云存储(S3, GCS) Zarr 3 通过URI字符串或`FsspecStore`使用**fsspec**后端(优先于旧版的`S3Map`/`GCSMap`)。 ```python import zarr # S3 — 优先使用IAM角色/配置文件;fsspec会自动处理提供商凭据发现。 # 切勿将凭据值打印、记录或复制到提示词或笔记本中。 z = zarr.create_array( store="s3://my-bucket/path/to/array.zarr", shape=(1000, 1000), chunks=(100, 100), dtype="f4", storage_options={"anon": False}, ) z[:] = data # GCS — 优先使用工作负载身份或gcloud应用默认凭据。 z = zarr.open_array( "gs://my-bucket/path/to/array.zarr", mode="r", storage_options={"project": "my-project"}, ) # 显式存储(任意fsspec文件系统) from zarr.storage import FsspecStore store = FsspecStore.from_url("s3://my-bucket/data.zarr", storage_options={"anon": False}) root = zarr.open_group(store=store, mode="r+") ``` 云端后端通过提供商 SDK/fsspec 后端读取凭据。不要检查大范围的`.env`文件;如果用户明确需要帮助调试认证问题,应要求提供脱敏后的配置,并只读取他们明确批准的具名提供商变量。将所有`import zarr`、`import dask`、`import h5py`、`import xarray`示例都当作第三方包导入,而不是随附的脚本文件。 **云存储最佳实践**: - 使用合并元数据减少延迟:`zarr.consolidate_metadata(store)` - 将分块大小与云对象大小对齐(通常5-100 MB最佳) - 使用Dask进行大规模数据的并行写入 - 考虑使用分片减少对象数量 ## 组和层次结构 组以层次结构组织多个数组,类似于目录或HDF5组。 ### 创建和使用组 ```python # 创建根组 root = zarr.group(store='data/hierarchy.zarr') # 创建子组 temperature = root.create_group('temperature') precipitation = root.create_group('precipitation') # 在组内创建数组 temp_array = temperature.create_array( name='t2m', shape=(365, 720, 1440), chunks=(1, 720, 1440), dtype='f4' ) precip_array = precipitation.create_array( name='prcp', shape=(365, 720, 1440), chunks=(1, 720, 1440), dtype='f4' ) # 使用路径访问 array = root['temperature/t2m'] # 可视化层次结构 print(root.tree()) # 输出: # / # ├── temperature # │ └── t2m (365, 720, 1440) f4 # └── precipitation # └── prcp (365, 720, 1440) f4 ``` ### 组API(v3) 使用`create_array` / `require_array`(h5py风格的`create_dataset` / `require_dataset`在v3中已移除): ```python root = zarr.group('data.zarr') arr = root.create_array('my_data', shape=(1000, 1000), chunks=(100, 100), dtype='f4') grp = root.require_group('subgroup') arr2 = grp.require_array('array', shape=(500, 500), chunks=(50, 50), dtype='i4') ``` ## 属性和元数据 使用属性将自定义元数据附加到数组和组: ```python # 向数组添加属性 z = zarr.zeros((1000, 1000), chunks=(100, 100)) z.attrs['description'] = 'Temperature data in Kelvin' z.attrs['units'] = 'K' z.attrs['created'] = '2024-01-15' z.attrs['processing_version'] = 2.1 # 属性以JSON形式存储 print(z.attrs['units']) # 输出:K # 向组添加属性 root = zarr.group('data.zarr') root.attrs['project'] = 'Climate Analysis' root.attrs['institution'] = 'Research Institute' # 属性与数组/组一起持久化 z2 = zarr.open('data.zarr') print(z2.attrs['description']) ``` **重要**:属性必须是JSON可序列化的(字符串、数字、列表、字典、布尔值、null)。 ## 与NumPy、Dask和Xarray的集成 ### NumPy集成 Zarr数组实现NumPy数组接口: ```python import numpy as np import zarr z = zarr.zeros((1000, 1000), chunks=(100, 100)) # 直接使用NumPy函数 result = np.sum(z, axis=0) # NumPy操作Zarr数组 mean = np.mean(z[:100, :100]) # 转换为NumPy数组 numpy_array = z[:] # 将整个数组加载到内存 ``` ### Dask集成 Dask对Zarr数组提供延迟、并行计算: ```python import dask.array as da import zarr # 创建大型Zarr数组 z = zarr.open('data.zarr', mode='w', shape=(100000, 100000), chunks=(1000, 1000), dtype='f4') # 加载为Dask数组(延迟,无数据加载) dask_array = da.from_zarr('data.zarr') # 执行计算(并行,核心外) result = dask_array.mean(axis=0).compute() # 并行计算 # 将Dask数组写入Zarr large_array = da.random.random((100000, 100000), chunks=(1000, 1000)) da.to_zarr(large_array, 'output.zarr') ``` **好处**: - 处理大于内存的数据集 - 跨分块自动并行计算 - 与分块存储的高效I/O ### Xarray集成 Xarray提供带有Zarr后端的标记多维数组: ```python import xarray as xr import zarr # 以Xarray Dataset形式打开Zarr存储(延迟加载) ds = xr.open_zarr('data.zarr') # 数据集包含坐标和元数据 print(ds) # 访问变量 temperature = ds['temperature'] # 执行标记操作 subset = ds.sel(time='2024-01', lat=slice(30, 60)) # 将Xarray Dataset写入Zarr ds.to_zarr('output.zarr') # 从头创建带有坐标 ds = xr.Dataset( { 'temperature': (['time', 'lat', 'lon'], data), 'precipitation': (['time', 'lat', 'lon'], data2) }, coords={ 'time': pd.date_range('2024-01-01', periods=365), 'lat': np.arange(-90, 91, 1), 'lon': np.arange(-180, 180, 1) } ) ds.to_zarr('climate_data.zarr') ``` **好处**: - 命名维度和坐标 - 基于标签的索引和选择 - 与pandas集成用于时间序列 - 气候/地理空间科学家熟悉的NetCDF类接口 ## 并行计算与线程安全 Zarr在内部使用异步I/O。针对远程存储或Dask密集型工作负载调整并发度: ```python import zarr # 更高的值可以提升远程吞吐量;当Dask已经提供了大量工作线程时, # 较低的值可以减少压力。 zarr.config.set({ "async.concurrency": 8, "threading.max_workers": 8, }) ``` 旧的`synchronizer`参数(`ThreadSynchronizer`、`ProcessSynchronizer`)**在Zarr-Python 3中尚不可用**。请改用以下模式: - **读取:** 跨线程/进程始终安全。 - **写入:** 当每个worker写入**互不重叠的分块**时是安全的;大多数存储都支持原子分块写入。 - **重叠写入:** 在同步器功能恢复前,需通过外部方式协调(文件锁、工作流设计等)。 对于Dask密集型工作负载,可将总并发I/O大致估算为`dask_threads × async.concurrency`,若存储或内存出现饱和,则调低Zarr的并发设置。 ## 合并元数据 对于有许多数组的层次存储,将元数据合并到单个文件中以减少I/O操作: ```python import zarr # 创建数组/组后 root = zarr.group('data.zarr') # ... 创建多个数组/组 ... # 合并元数据 zarr.consolidate_metadata('data.zarr') # 使用合并的元数据打开(更快,尤其是在云存储上) root = zarr.open_consolidated('data.zarr') ``` **好处**: - 将元数据读取操作从N(每个数组一个)减少到1 - 对云存储至关重要(减少延迟) - 加快`tree()`操作和组遍历 **注意**: - 如果数组更新而不重新合并,元数据可能会过时 - 不适合频繁更新的数据集 - 多写入器场景可能有不一致的读取 ## 性能优化 ### 最佳性能清单 1. **分块大小**:目标为每个分块1-10 MB ```python # 对于float32:1MB = 262,144个元素 chunks = (512, 512) # 512×512×4字节 = ~1MB ``` 2. **分块形状**:与访问模式对齐 ```python # 行式访问 → 分块跨越列:(小, 大) # 列式访问 → 分块跨越行:(大, 小) # 随机访问 → 平衡:(中, 中) ``` 3. **压缩**:根据工作负载选择 ```python # 交互式/快速:BloscCodec(cname='lz4') # 平衡:BloscCodec(cname='zstd', clevel=5) # 最大压缩:GzipCodec(level=9) ``` 4. **存储后端**:与环境匹配 ```python # 本地:LocalStore(默认) # 云:fsspec URI 或 FsspecStore + 合并元数据 # 临时:MemoryStore ``` 5. **分片**:用于大规模数据集 ```python # 当有数百万个小分块时 shards=(10*chunk_size, 10*chunk_size) ``` 6. **并行I/O**:大型操作使用Dask ```python import dask.array as da dask_array = da.from_zarr('data.zarr') result = dask_array.compute(scheduler='threads', num_workers=8) ``` ### 分析和调试 ```python # 打印详细的数组信息 print(z.info) # 输出包括: # - 类型、形状、分块、数据类型 # - 序列化器和压缩器 # - 存储大小(压缩与未压缩) # - 存储位置 # 检查存储大小 print(f"Compressed size: {z.nbytes_stored / 1e6:.2f} MB") print(f"Uncompressed size: {z.nbytes / 1e6:.2f} MB") print(f"Compression ratio: {z.nbytes / z.nbytes_stored:.2f}x") ``` ## 常见模式和最佳实践 ### 模式:时间序列数据 ```python # 以时间为第一维度存储时间序列 # 这允许有效追加新时间步 z = zarr.open('timeseries.zarr', mode='a', shape=(0, 720, 1440), # 从0个时间步开始 chunks=(1, 720, 1440), # 每个分块一个时间步 dtype='f4') # 追加新时间步 new_data = np.random.random((1, 720, 1440)) z.append(new_data, axis=0) ``` ### 模式:大型矩阵操作 ```python import dask.array as da # 在Zarr中创建大型矩阵 z = zarr.open('matrix.zarr', mode='w', shape=(100000, 100000), chunks=(1000, 1000), dtype='f8') # 使用Dask进行并行计算 dask_z = da.from_zarr('matrix.zarr') result = (dask_z @ dask_z.T).compute() # 并行矩阵乘法 ``` ### 模式:云原生工作流 ```python import zarr path = "s3://my-bucket/data.zarr" z = zarr.create_array( store=path, shape=(10000, 10000), chunks=(500, 500), dtype="f4", storage_options={"anon": False}, ) z[:] = data zarr.consolidate_metadata(path) z_read = zarr.open_consolidated(path, storage_options={"anon": False}) subset = z_read[0:100, 0:100] ``` ### 模式:格式转换 ```python # HDF5到Zarr import h5py import zarr with h5py.File('data.h5', 'r') as h5: dataset = h5['dataset_name'] z = zarr.array(dataset[:], chunks=(1000, 1000), store='data.zarr') # NumPy到Zarr import numpy as np data = np.load('data.npy') z = zarr.array(data, chunks='auto', store='data.zarr') # Zarr到NetCDF(通过Xarray) import xarray as xr ds = xr.open_zarr('data.zarr') ds.to_netcdf('data.nc') ``` ## 常见问题和解决方案 ### 问题:性能缓慢 **诊断**:检查分块大小和对齐 ```python print(z.chunks) # 分块大小是否合适? print(z.info) # 检查压缩比 ``` **解决方案**: - 将分块大小增加到1-10 MB - 使分块与访问模式对齐 - 尝试不同的压缩编解码器 - 使用Dask进行并行操作 ### 问题:内存使用高 **原因**:将整个数组或大型分块加载到内存 **解决方案**: ```python # 不要加载整个数组 # 错误:data = z[:] # 正确:分块处理 for i in range(0, z.shape[0], 1000): chunk = z[i:i+1000, :] process(chunk) # 或使用Dask进行自动分块 import dask.array as da dask_z = da.from_zarr('data.zarr') result = dask_z.mean().compute() # 分块处理 ``` ### 问题:云存储延迟 **解决方案**: ```python # 1. 合并元数据 zarr.consolidate_metadata(store) z = zarr.open_consolidated(store) # 2. 使用适当的分块大小(云为5-100 MB) chunks = (2000, 2000) # 云用更大的分块 # 3. 启用分片 shards = (10000, 10000) # 组许多分块 ``` ### 问题:并发写入冲突 **解决方案**:设计工作流,让每个进程/线程只写入各自独立的分块。Zarr-Python 3 尚不支持`ThreadSynchronizer` / `ProcessSynchronizer`;参见`references/v3_migration.md`。 ## 其他资源 ### 内置参考资料 | 文件 | 内容 | |------|------| | `references/api_reference.md` | 函数签名、存储、编解码器、索引 | | `references/v3_migration.md` | Zarr-Python 2→3 的破坏性变更及进行中的功能 | ### 官方上游资源 - **文档**:https://zarr.readthedocs.io/en/stable/ - **3.0 迁移指南**:https://zarr.readthedocs.io/en/stable/user-guide/v3_migration/ - **存储后端**:https://zarr.readthedocs.io/en/stable/user-guide/storage/ - **Zarr 规范**:https://zarr-specs.readthedocs.io/ - **GitHub**:https://github.com/zarr-developers/zarr-python - **开发者聊天室**:https://ossci.zulipchat.com/#narrow/channel/423692-Zarr-Python **相关库:** [Xarray](https://docs.xarray.dev/)、[Dask](https://docs.dask.org/)、[NumCodecs](https://numcodecs.readthedocs.io/)