1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
//! High-performance sequential dataset loading, in the
//! [WebDataset](https://github.com/webdataset/webdataset) format.
//!
//! A WebDataset is a set of tar archives ("shards"). Inside a shard, the files
//! that share a basename make up one training sample:
//!
//! ```text
//! imagenet-000000.tar
//! n03991062_24866.jpg n03991062_24866.cls
//! n03995372_9042.jpg n03995372_9042.cls
//! ```
//!
//! Reading is purely sequential, which is what makes the format fast: no seeks,
//! no index, no per-file requests. A shard reads at the full bandwidth of the
//! device or the network link, and a dataset is just a list of URLs — nothing
//! needs to be registered, converted, or mounted first.
//!
//! # Reading a dataset
//!
//! ```no_run
//! use webdataset::{Decoder, WebDataset};
//! use webdataset::filters::{SampleIteratorExt, TupleIteratorExt};
//!
//! let dataset = WebDataset::builder("https://host/imagenet-{000000..000146}.tar")
//! .shard_shuffle(100)
//! .build()?
//! .shuffle(1000)
//! .decode(Decoder::default());
//!
//! for batch in dataset.iter().to_tuple(["jpg;png", "cls"]).batched(64, true) {
//! let batch = batch?;
//! let images = &batch[0];
//! let labels = &batch[1];
//! # let _ = (images, labels);
//! }
//! # Ok::<(), webdataset_core::Error>(())
//! ```
//!
//! # Writing a dataset
//!
//! ```
//! use std::sync::Arc;
//! use webdataset::encode::DefaultEncoder;
//! use webdataset::{Sample, ShardWriter, Value};
//!
//! # let dir = tempfile::tempdir()?;
//! # let pattern = dir.path().join("train-%06d.tar");
//! # let pattern = pattern.to_str().unwrap();
//! let mut writer = ShardWriter::new(pattern)?
//! .with_encoder(Arc::new(DefaultEncoder::new()))
//! .with_max_count(10_000);
//!
//! for i in 0..100 {
//! let mut sample = Sample::with_key(format!("sample{i:06}"));
//! sample.insert("cls", Value::Int(i % 10));
//! sample.insert("txt", Value::Text(format!("sample number {i}")));
//! writer.write(&sample)?;
//! }
//! writer.close()?;
//! # Ok::<(), webdataset_core::Error>(())
//! ```
//!
//! # How a pipeline is put together
//!
//! [`WebDataset`] assembles the usual stages for you, but the pieces are public
//! and a pipeline can be built by hand:
//!
//! ```text
//! SimpleShardList a list of shard URLs
//! -> SplitByNode keep this rank's shards
//! -> SplitByWorker keep this worker's shards
//! -> Shuffle shuffle the shard order
//! -> ShardsToSamples open each shard, group its files into samples
//! -> Shuffle shuffle samples within a buffer
//! -> Decode turn bytes into images, tensors, JSON
//! ```
//!
//! Everything up to this point is a [`Stage`] producing
//! [`Sample`]s, so it can be re-run each epoch. Transformations that change the
//! item type — projecting to tuples, batching — are ordinary iterator adapters
//! from [`filters`], applied to [`WebDataset::iter()`].
//!
//! # Scaling out
//!
//! [`DataLoader`] runs one copy of the pipeline per worker thread. Shards are
//! divided between workers by the [`SplitByWorker`] stage, and between
//! distributed processes by [`SplitByNode`]. When exact partitioning is
//! awkward — many nodes, few shards — use `resampled(true)` instead and let
//! each worker draw shards with replacement.
//!
//! # Errors
//!
//! Every stream yields `Result<Sample>`. What happens when a shard is corrupt
//! or a field fails to decode is decided by a
//! [`Handler`]: forward the error, drop the sample,
//! or end the stream. Streaming a petabyte means meeting some bad bytes, so
//! `warn_and_continue` is a common choice.
//!
//! # Features
//!
//! | feature | adds |
//! |---|---|
//! | `threads` *(default)* | per-shard read-ahead and the multi-worker loader |
//! | `subprocess` *(default)* | `pipe:` and the `curl`/`gsutil`/`ais` schemes |
//! | `yaml` *(default)* | multi-source dataset specifications |
//! | `image` | `.jpg`, `.png` and friends, via the `image` crate |
//! | `msgpack` | `.mp` and `.msg` |
//! | `cbor` | `.cbor` |
//! | `npz` | NumPy `.npz` archives |
//! | `zstd`, `bzip2`, `xz` | shards in those containers |
//! | `async` | read shards from any `AsyncRead`, yielding a `Stream` |
//! | `wasm-js` | host randomness on `wasm32-unknown-unknown` |
//! | `full` | every format, plus threads, subprocesses and async |
//!
//! `.npy`, `.ten`, `.json`, `.txt`, `.cls` and gzip need no features.
//!
//! # Async
//!
//! With the `async` feature, [`asynch`] mirrors everything above: the same
//! archive parser, the same decoders, the same shuffling and batching, driven
//! by futures and producing a [`Stream`](futures_core::Stream) rather than an
//! [`Iterator`]. Reach for it when shards arrive over a network — a blocking
//! reader holds a thread for the whole of a transfer, an async one holds only a
//! task, and raising
//! [`concurrency`](asynch::AsyncWebDatasetBuilder::concurrency) overlaps the
//! fetches.
//!
//! Only the I/O is async; decoding and batching are the same code the blocking
//! pipeline runs, which is why the two read identically — a property the test
//! suite checks shard by shard. See the [`asynch`] module for a worked example.
//!
//! # WebAssembly
//!
//! WebAssembly has no threads to spawn and no processes to run, so turn both
//! features off and hand the shard bytes over yourself with
//! [`MemoryOpener`]. Everything above the transport — the archive parser, the
//! decoders, shuffling, batching — is unchanged.
//!
//! ```
//! use std::sync::Arc;
//! use webdataset::pipeline::DataPipeline;
//! use webdataset::shardlists::SimpleShardList;
//! use webdataset::sources::{MemoryOpener, ShardsToSamples};
//! use webdataset::stages::Decode;
//!
//! # let bytes: Vec<u8> = std::fs::read("../../testdata/sample.tgz")?;
//! let opener = MemoryOpener::new().with("mem://shard-000.tar", bytes);
//! let dataset = DataPipeline::new()
//! .with(SimpleShardList::verbatim(["mem://shard-000.tar"]))
//! .with(ShardsToSamples::new(Arc::new(opener)))
//! .with(Decode::basic());
//!
//! assert_eq!(dataset.iter().count(), 90);
//! # Ok::<(), webdataset_core::Error>(())
//! ```
//!
//! For any other transport, register a scheme with
//! [`webdataset_io::register_scheme`].
//!
//! `webdataset-core` and `webdataset-tenbin` go further and build without the
//! standard library at all; see their documentation.
pub use ;
pub use ;
pub use ;
pub use DefaultEncoder;
pub use ;
pub use DataLoader;
pub use ;
pub use ;
pub use ;
pub use ;
pub use ;
/// The core data model: samples, values, tensors, and errors.
pub use webdataset_core as core;
/// URL opening and shard caching.
pub use webdataset_io as io;
/// Streaming tar readers and writers.
pub use webdataset_shard as shard;
/// The `.ten` binary tensor format.
pub use webdataset_tenbin as tenbin;
pub use ;
pub use ;
pub use ;