Skip to content
ccrawl

Running the recrawl fleet

Install ccrawl on three machines, run a shard of the work list on each under systemd, publish as it goes, and know what to do when one falls behind.

A recrawl of the published domain list is 121 million rows, and of the URL list about 2.1 billion. At any rate one machine can hold, that is months of running, which makes this an operations problem rather than a command to type. This page is the runbook: what gets installed, how it is started and stopped, what happens across a reboot, and what to do when one of the three machines stops keeping up.

The engine itself is described in building a recrawl engine. This page assumes it and talks only about running it on more than one box.

The shape

Each machine takes one shard of the work list and publishes into a shared dataset repo. Two units run per machine per work list, both templated on the name of the list:

Unit Does
ccrawl-recrawl@domains Fetches this machine's third of the domain list and writes Parquet shards to local disk
ccrawl-publish@domains Watches that directory and commits each shard to the hub as it closes, then deletes it
ccrawl-recrawl.target One handle for everything this machine crawls

The publisher runs beside the crawl and not after it. A run of this length that published at the end would need months of disk, and none of these machines has that. It also means the dataset is readable from the first hour rather than the last day.

The partition key is the registered domain, so a site and everything under it stays on one machine. That is what keeps the politeness clock meaningful: two machines crawling the same site would each think they were spacing their requests properly, and the site would see twice the rate.

Installing

From a checkout, on a machine that can ssh to all three:

deploy/install.sh                                  # all three, domain list
deploy/install.sh --servers server1 --kind urls    # one machine, url list
deploy/install.sh --start                          # install and turn it on

It builds one binary, copies it to each server, checks the hash on the far side before moving it into place, installs the units, and enables them. It does not start anything unless asked, because starting a crawl is a decision about somebody else's bandwidth.

The reason it builds rather than copying whatever is on the box is the state the fleet was found in. server3 was running 0.5.0 from July, and the other two were dev builds from two different afternoons. A fleet running three binaries is a fleet where a rate difference between two machines means nothing, so the script prints the version off all three when it finishes and that output is the thing to read.

Two files are left alone by every install after the first:

/etc/ccrawl/recrawl-<kind>.env holds the shard number and the tuning. It is written once, from the example in deploy/env, with the shard and the server name filled in. A deploy that reset it would silently undo whatever the last person measured.

/etc/ccrawl/hf.env holds the token, mode 0600, read by the publisher and not by the crawl. It is created empty, so the first commit fails loudly rather than the hundredth failing quietly. Put the token in before starting a publisher.

Starting and stopping

systemctl start ccrawl-recrawl.target     # everything this machine crawls
systemctl stop ccrawl-recrawl.target      # all of it, cleanly
systemctl status ccrawl-recrawl@domains   # one crawl
journalctl -fu ccrawl-recrawl@domains     # watch it

A stop is safe at any moment. The checkpoint is written at a batch boundary and holds a part number and a row offset, so a run that is stopped and started again refetches the batch it was in and whatever the pool had in the air behind it, which is a few hundred rows. It never skips. The unit allows ninety seconds to shut down, which is far more than a batch takes to drain.

A reboot needs nothing done to it. The target is enabled, the units come back with the machine, and each one reads its checkpoint and carries on. This is worth testing once per machine rather than finding out during one.

What the numbers in the env file mean

CCRAWL_WORKERS is the width of the pool, and it is the only one of these that moves the rate.

The rate is workers times yield over time per page, and on the domain corpus the yield is about 55 percent. Everything else in this file has been measured against the live list and none of it changed the rate, which is the section below. Workers is also what decides the memory, so read the memory section before raising it.

CCRAWL_WRITERS is how many output files are open at once, each with its own encoder. Raise it only when a run reports the writer busy in the high eighties. At four writers the share drops to about a third each, and the page rate moves by a few percent, so it is headroom rather than speed.

--extractors is how wide the extract pool is, and it is a separate number from the workers on purpose. Fetching waits on the far end, so it wants to be hundreds wide and costs almost nothing while it waits, and rendering waits on nothing at all. The default is four times the core count, which is not the number you would guess, and the section below has the runs it came from. There is no env variable for it because the default is right on both machines and the flag is for the day it is not.

CCRAWL_DNS_LOOKUPS bounds the lookups in flight. Every row of a domain corpus is a name nobody has looked up yet, so DNS is a per page cost there rather than a detail of the fetch. The end of a run prints the peak against the bound. A peak well under the bound means DNS is not the constraint and raising it will do nothing. A peak sitting on the bound does not mean the opposite, which is the trap and is measured below.

CCRAWL_SHARD_SIZE is how much payload goes into a shard before it is sealed, counted uncompressed. It decides how many files land on the hub and how much work a crash replays, and it does not decide the memory even though a shard buffers in memory until it is sealed.

CCRAWL_TIMEOUT is the budget for one fetch and its retries together. Leave it at 30 seconds, for the reason below.

Three dials that look like the ceiling and are not

Each of these was measured on server3 against the live domain list, 20000 rows a run from a fresh row offset, runs alternated so a drift in the box load falls on both settings rather than on one. They are written down because all three are things the summary line invites you to turn, and turning any of them wastes a day.

The DNS bound. The peak pegs at whatever the bound is set to, so pegging proves there is demand and not that the queue costs anything.

dns-lookups pages a second seconds an item peak idle
32 19.2 5.498 32 27%
128 19.3 5.289 128 29%
512 20.7 4.986 216 28%
512 23.9 4.144 200 31%

Four times the bound moved the rate by half a percent. The two runs at 512 differ from each other by 15 percent, which is the box, so even the gain at 512 is inside the noise.

The fetch timeout. Roughly 1250 rows in every 20000 time out at 30 seconds, which looks like a quarter of the pool held for the whole run and returning nothing. Cutting it does make each item faster and it does not make the run faster.

timeout pages a second seconds an item timeouts per 20000
30s 29.8 4.069 1246
10s 28.1 4.418 1438
5s 25.9 3.828 3779
5s 28.8 3.673 3170
10s 27.1 4.529 1573
30s 25.6 4.759 1246

At five seconds the time per item drops about 15 percent and the timeouts nearly triple, so the yield drops about as much as the item got faster and the two cancel. The averages are 27.7, 27.6 and 27.4 pages a second, and the spread inside the 30 second setting alone is 25.6 to 29.8. A short timeout does not skip slow pages, it converts them into failures.

The shard size. A shard buffers in memory until it is sealed, so a 512 MB target across four writers looks like 2 GB of buffer before a worker holds anything. Measured on server2 at 48 workers and 2 writers, alternated:

shard size peak resident
20 MB 1.042 GB
512 MB 1.177 GB
512 MB 1.190 GB
20 MB 1.299 GB

The smallest setting produced the highest peak of the four. The spread inside a setting is 25 percent and the difference between settings is 1 percent, so the rows are compressed as they are buffered and the buffer is not where the memory goes. The memory is on the worker side, which is the next section.

The summary at the end of a run is written to be read in this order. The failure breakdown says what the corpus is costing: on the domain list about 5500 rows in 20000 are names that do not resolve, 1400 are broken TLS handshakes and 1200 time out, and none of that is the machine's fault or something a wider pool fixes. The timing line says what the machine is costing: how much of the run the pool spent idle, and how much of it each writer was busy. The feeder line under it says what the work list is costing, and it is there because the timing line can say the pool was idle and cannot say why. There is one feeder and it has two states, reading the next row off the work list or waiting for a worker to take the row it is holding, so those are shares of a single goroutine's clock rather than of the pool and they should not be added to the percentages above them. The line leads with how long the feeder was alive because the rest of it is a share of that and not of the run. A feeder stops as soon as the work list ends or --max-pages fires, and the run carries on until the pool drains, and on a bounded probe that tail can be most of the run: measured on server2 with a 400 page limit, the feeder finished in 22 seconds of a 62 second run because one straggler sat out the whole 30 second timeout at the end. Reading large is a work list the run cannot get through fast enough to fill the pool, which is a Parquet and bandwidth problem. Waiting for a free worker large is the pool being the constraint, which is the healthy reading and the one where the phase shares are worth acting on. The rows a second at the end is the ceiling on the whole run, because no run fetches more pages a second than its feeder hands out rows.

The URL list is sorted by host, and that is its ceiling

The two work lists need different things from the same crawler, and the difference is the order they are in.

The domain list is one row per host, so a batch of it is two thousand different sites and every worker has something to fetch. The URL list comes out of Common Crawl's index in SURT order, which is sorted by host and then by path, so a batch of it is one site. The politeness delay gives a host one request per second no matter how many workers are free, so a pool reading that list in order runs at one page a second and the rest of the pool waits.

Measured on server3 at 32 workers against the live URL list, before the fix: 596 pages in ten minutes, and all 596 of them from vkbn.ru.

The crawler now reads the URL list through a reorder buffer that keeps reading until it holds one host per worker, then hands rows out one host at a time. Each site keeps its own pages in work list order, so this is a rotation and not a shuffle, and two runs over the same list hand out the same sequence. The domain list is unaffected, because a read of it already satisfies the buffer on the first batch.

The cost is replay after a crash. The checkpoint can only name the oldest row that has not finished, so rows held back in the buffer hold the checkpoint back with them. The buffer caps how many rows it holds at sixteen batches, and that is not the same as capping how far back the checkpoint sits, which is worth being plain about. A big site drains at one page a second whatever the pool is doing, so its rows stay in hand while the reader runs on through everybody else. Measured on the live URL run on server3, the oldest row in hand was 126287 while the reader was at 489026, and that gap grows for as long as that site has rows left. Nothing is lost by it, since the pages in between are fetched and written and committed as usual, but a kill refetches them. The way out is a work list that is not sorted by host, which is a change to how the URL list is published rather than to how it is read.

Measured on server3 against the live URL list, with the domain run going on the same box:

workers distinct hosts in a two minute window fetched pages a second
32, before the buffer 1 1.0
32 80 6.3
96 17.6
256 23.3

The rate follows the worker count because the worker count is what the buffer reads ahead for, and each worker holding its own host is what turns the one request per second per host into one request per second per worker. It stops following it somewhere before 256, where the box runs out of whatever it runs out of first, and 256 workers on the URL list also halved the domain run beside it. 96 is where server3 is left, as the most rate per unit of memory rather than the most rate.

Three stages, three widths

A recrawl does three jobs and they answer to three different limits. Fetching waits on the far end and wants to be hundreds wide. Rendering a page into text and Markdown waits on nothing at all and wants to be about as wide as there are cores, because past that it is only competing with itself. Encoding and compressing Parquet is a third again.

For a long time the first two shared a pool, and the note at the top of recrawl_extract.go makes the case for that: the run is network bound, a worker is waiting on a socket most of the time, and rendering fits in the gap the fetches leave. That was true when it was written and it stopped being true as the pool got wider. Measured on server2 at 96 workers, a run with rendering inline reported fetching at 44 percent of the pool and extracting at 27, which is 27 percent of ninety six workers spent on a job six cores can do. The pool was not wide, it was blocked, and there was no way to tell the two widths apart because there was only one of them.

So a fetch worker now hands the response over and goes back to the network, and a pool sized from the core count renders. --extractors is the width of that pool and the summary gains a line about it:

recrawl run: extract pool: 6 wide, rendering 71% of it, waiting for the sink 4%, idle 25%

Read it the same way as the writer share and not the same way as the worker line, because it is a share of its own width times the wall clock rather than of the fetch pool's. Rendering near the whole pool is an extract pool that is the constraint and wants more of the machine. Waiting for the sink large is an extract pool that is keeping up and queueing behind the writers, which is what CCRAWL_WRITERS is for and not what a wider extract pool is for. Both small is a pool with room to spare, and then the worker line is the one to read.

The worker line no longer has an extracting share, for the same reason. It used to, when rendering was worker time, and leaving it there would report fetch workers busy with work no fetch worker did.

How wide the extract pool wants to be

The core count is the obvious answer and it is wrong, by a lot.

Measured on server2, 96 workers, 3 writers, 15000 pages of the URL list with extraction on. Two sweeps in two different windows, each one bracketed by the same width run twice, since this box's load moves on its own and only the ordering inside a sweep is comparable:

extractors rate rendering waiting for the sink
6, the core count 32.5 then 31.2 86% 10%
12 35.2 74% 21%
24 40.3 65% 26%
extractors rate rendering waiting for the sink
24 30.9 then 30.3 68% then 65% 27% then 32%
48 33.3 54% 37%
96 33.8 39% 46%

The second window is slower than the first at the same width, which is what the brackets are there to show, and in both of them wider is faster.

The reason is that these are shared machines. Our runnable goroutines are what win a share of the cores against the other tenants, and a pool of six on a box with a load average of twelve gets a sixth of the machine however much work it is holding. A pool sized from nproc is sized for a machine we do not have.

It flattens between 48 and 96, and the reason it flattens is in the same two lines: rendering falls to 39 percent while waiting for the sink climbs to 46 and the writers reach 61 percent busy. Past that width the extract pool is not the constraint any more and the answer is CCRAWL_WRITERS. Four times the cores is 24 on server2 and 32 on server3, which is near where that handover happens on both.

There are two queues now rather than one, so there are two byte bounds, and --queue-mb is split between them:

recrawl run: extract queue: 67.1 MB bound, 12.4 MB at its highest, 18% of the bound, write queue: 67.1 MB bound, 20.1 MB at its highest, 30% of the bound

The split is why the figure is halved rather than each stage getting the whole number. What --queue-mb promises is what the run holds between the fetches and the disk, and a run holding that much in two places would cost twice what it said. A run that is not rendering has one queue and it gets the whole budget.

Memory is the binding constraint, not CPU

This is the thing to know before tuning anything.

server1 server2 server3
CPU 4 6 8
RAM 5 GB 11 GB 23 GB
Available RAM under 1 GB about 1 GB about 3 GB
Free disk 156 GB 20 GB 27 GB

These are shared machines with other work on them, and the available column is what is actually free rather than what is installed.

Measured peaks, all on the live domain list:

workers writers shard size peak resident
48 2 20 MB to 512 MB 1.0 to 1.3 GB
64 2 20 MB 1.1 GB
256 1 512 MB 2.7 GB
256 4 512 MB 4.5 GB

Neither of the 256 worker numbers fits in what server1 has free today.

The rule of thumb between the two ends is roughly 15 to 20 MB of resident memory per worker, over a floor of about a gigabyte that does not move much below 64 workers. It is a rule of thumb and not a formula: it is fitted to a handful of points on one corpus, and the page sizes on a different slice of the work list would move it. Use it to pick a starting width and then watch the run rather than trusting the arithmetic.

The queue between the pool and the sink is the part of that memory a run can say a number about, and --queue-mb is the number. It defaults to 128 MB, which is about four hundred average pages of this corpus and more slack than a pool of any width needs to ride out a burst. It used to be a count of items with a comment beside it saying that a page averages 300 KB, which is a way of admitting that what it meant to bound was bytes, and a count cannot do that when a body runs from a kilobyte to the 10 MB cap. The run that made the case for changing it reached 3.6 GB resident on server2 at 288 workers and 6 writers, on a box with 2 GB free, and was killed along with two of another tenant's processes. Nothing in the settings had said that would happen because the settings were counting the wrong thing.

The queue line at the end of a run reports the bound and the most that was ever held against it, and the two are only useful together. A peak well under the bound means the queue was never the reason a worker waited, and the pool can be given more workers. A peak sitting on the bound means the sink is the ceiling and the workers are queueing behind memory rather than behind the network, which is what --writers is for and not what a larger queue is for. Raising the bound in that second case buys a longer wait and the same rate.

When the queue line and the timing line have both been read and the answer is still a phase rather than a knob, --profile-dir and --profile-name write a CPU profile and a heap profile for the run. Give the two runs being compared different names and diff them with go tool pprof -base first.cpu.pprof second.cpu.pprof, because a change is worth judging on the difference it made and not on either profile alone. The heap profile is the one to reach for when a run is using more memory than the queue bound explains, which is most of the time: the queue is bounded and the encoder buffers are not, and a shard buffers until it is sealed. Profiling is off unless asked for and no unit sets it, since it costs a few percent and writes tens of megabytes.

The first thing that profile found is worth knowing before setting CCRAWL_WRITERS. A heap profile of a 96 worker run showed 750 MB in use with 54 percent of it inside zstd, against 124 MB of Parquet column buffers and 136 MB of response bodies. zstd holds that state per concurrent encoder, every writer used to open four of them, so the encoders were writers times four rather than four in total, and at six writers that is twenty four of them holding somewhere over a gigabyte before a single page is in a buffer. The encoders are divided across the writers now, so the product stays where it was at one writer and raising CCRAWL_WRITERS costs the column buffers and not the compressor as well. The memory table further down was measured before that change and is the pessimistic reading of what a wide run costs.

The second thing the profile found was a row count in the same job a byte bound should have had, one layer further in. A writer holds a batch of rows before it hands them to the encoder, that batch used to be 500 rows and nothing else, and 500 rows of the URL corpus is about 72 MB, so three writers sat on 217 MB of bodies waiting to reach a count. Bodies carried their own slack on top of that, because a buffer that grows by doubling keeps an array up to twice the length of what it read and that array travels into the row, the queue and the batch. A batch now hands over on whichever of the count or 16 MB fills first, a body is sized from Content-Length where the far end gives one, and a body left holding real slack is copied to its own length before it goes anywhere. On server2 at 96 workers and 3 writers that took bytes.growSlice from 222 MB to 1 MB, in use heap from 493 MB to 136 MB, and peak resident from 1.69 GB to 0.85 GB.

So the width has to be set against the memory the box actually has at the moment, and that means checking free -g before raising CCRAWL_WORKERS rather than copying a number from another machine. install.sh writes a MemoryHigh of 70 percent of installed RAM into a drop-in per machine. MemoryHigh throttles rather than kills, which is the right end for a crawler: a run that slows down is one that carries on, and a run that is OOM killed repeats its batch and may do it again.

The logs are the other thing on the disk, and they were the surprise. The run writes a JSON line per page to stdout, which is right for somebody watching a crawl and wrong for a unit that is up for months: an hour of it on server3 is about 390 MB of journal, so a machine would spend half a gigabyte a day recording lines nobody reads. The crawl unit sends stdout to /dev/null for that reason. Nothing an operator uses is lost, because the release the run picked, the failure breakdown, the timing line and every error are on stderr, and journalctl -u ccrawl-recrawl@domains still answers the questions it is asked. Run the binary by hand when a per page trace is what you want.

Disk is the constraint people expect and it is currently the one that is fine, because the publisher deletes each shard after committing it. The number to watch is not the free space, it is whether the free space is flat. A capture directory that grows is a publisher that has stopped, and that is the failure below.

When one machine falls behind

The three machines have different CPU, different memory and different neighbours, so they will not run at the same rate and are not meant to. Falling behind matters when it is a fault rather than a difference.

The publisher stopped and the crawl did not. Free disk drops steadily and the capture directory fills. journalctl -u ccrawl-publish@domains will usually say the hub rejected a commit or the token is wrong. The crawl can be left running while this is fixed if there is disk to spare, since the publisher picks up every closed shard it finds when it comes back. Shards are named by a hash of their contents, so a shard that was uploaded and not deleted is republished as a no-op rather than a duplicate.

The crawl is restarting in a loop. systemctl status shows a recent start and a low uptime, repeatedly. The unit gives up after ten starts in ten minutes and stays down, which is deliberate: at that point it is the binary or the config rather than the network, and a machine that sits still is one that gets noticed.

A restart is not free, and how much it costs depends on how deep into a part the checkpoint sits. Coming back means positioning the reader at the checkpoint row before a single page is fetched. That used to mean reading every row before it and throwing it away, which over ranged HTTP pulls every byte before it across the network. It seeks now, and the difference is what the depth is worth. Measured on server2 against part 0 of the live URL list, which holds 6890137 rows, timing a one page run from a fixed row and alternating the two binaries so drift on a loaded box cancels:

resume row before after
540000 13s, 9s 7s, 6s
6000000 39s, 43s 12s, 6s

So it is worth a few seconds where the fleet sits today and most of a minute near the end of a part, and it grows with depth rather than staying put. This matters more than it looks, because a machine that is restarting in a loop pays it on every loop.

The crawl is running and the rate is much lower than the other two. Read the summary lines rather than the rate. A writer share in the high eighties is the sink, and CCRAWL_WRITERS is the answer. A DNS peak well under the bound rules the resolver out. A high idle share with a low writer share is the pool waiting on something outside the machine, and on these machines that is usually the network or the neighbours rather than a setting, so read the load average before changing anything. Do not reach for CCRAWL_DNS_LOOKUPS or CCRAWL_TIMEOUT, which were both measured and neither moved the rate.

A machine has to be taken out of the fleet. Stop its target, then leave it. Do not repartition the remaining two, because the shard number is what decides which rows a machine owns, and changing CCRAWL_SHARDS from three to two moves every row on every machine and invalidates all three checkpoints. The right move is to fix the machine and let it resume, or to accept that a third of the work list is paused.

Where the output goes

open-index/ccrawl-recrawl-domains and open-index/ccrawl-recrawl-urls, one repo each, because the two lists finish on completely different schedules and nobody wants a card that averages them.

Each machine writes exactly one ledger file, ledger/<server>-shard<i>of<n>.csv, and never touches another machine's. Three machines committing at the same moment therefore cannot lose each other's numbers. The dataset card is generated from the union of every ledger on the hub, so it corrects itself on the next commit from any machine, and a machine that was down for a day does not leave a permanently wrong card behind.

Every row carries the extracted text and Markdown as well as the body, rendered as the page was fetched, so a published shard is usable without a second pass over the corpus.