This commit is contained in:
@@ -23,6 +23,7 @@ import { revisions as january2020Revisions } from '../scripts/upgrade-2020-01.mj
|
||||
import { revisions as february2020Revisions } from '../scripts/upgrade-2020-02.mjs';
|
||||
import { revisions as march2020Revisions } from '../scripts/upgrade-2020-03.mjs';
|
||||
import { revisions as april2020Revisions } from '../scripts/upgrade-2020-04.mjs';
|
||||
import { revisions as may2020Revisions } from '../scripts/upgrade-2020-05.mjs';
|
||||
|
||||
// This layer replaces archived source entries without losing their stable slug and date.
|
||||
export const editorialRevisions = [
|
||||
@@ -51,4 +52,5 @@ export const editorialRevisions = [
|
||||
...february2020Revisions,
|
||||
...march2020Revisions,
|
||||
...april2020Revisions,
|
||||
...may2020Revisions,
|
||||
];
|
||||
|
||||
@@ -0,0 +1,53 @@
|
||||
<svg xmlns="http://www.w3.org/2000/svg" viewBox="0 0 720 1210" role="img" aria-labelledby="title desc">
|
||||
<title id="title">Маршрут диагностики фоновой задачи</title>
|
||||
<desc id="desc">Вертикальная схема проводит расследование по одному jobId через запись задачи, outbox и публикацию, журнал worker, сохранённый результат либо карантинную очередь.</desc>
|
||||
<defs>
|
||||
<marker id="arrow" viewBox="0 0 12 12" refX="6" refY="6" markerWidth="10" markerHeight="10" orient="auto-start-reverse"><path d="M 1 1 L 11 6 L 1 11 z" fill="#516274" /></marker>
|
||||
<style>
|
||||
.title { fill: #182230; font: 700 34px system-ui, -apple-system, sans-serif; }
|
||||
.lead { fill: #516274; font: 400 22px system-ui, -apple-system, sans-serif; }
|
||||
.number { fill: #ffffff; font: 700 25px system-ui, -apple-system, sans-serif; }
|
||||
.step { fill: #182230; font: 700 28px system-ui, -apple-system, sans-serif; }
|
||||
.copy { fill: #354458; font: 400 23px system-ui, -apple-system, sans-serif; }
|
||||
.small { fill: #516274; font: 400 21px system-ui, -apple-system, sans-serif; }
|
||||
</style>
|
||||
</defs>
|
||||
<rect width="720" height="1210" fill="#f6f8fb" rx="28" />
|
||||
<text class="title" x="50" y="72">Диагностика по jobId</text>
|
||||
<text class="lead" x="50" y="108">Не перезапускать effect, пока не известна граница</text>
|
||||
|
||||
<rect x="48" y="152" width="624" height="142" rx="22" fill="#e7f0ff" stroke="#4a78bd" stroke-width="3" />
|
||||
<circle cx="100" cy="214" r="30" fill="#3167b1" />
|
||||
<text class="number" x="93" y="223">1</text>
|
||||
<text class="step" x="158" y="206">Запись задачи</text>
|
||||
<text class="copy" x="158" y="247">state, attempt, lastError, resultKey</text>
|
||||
<line x1="360" y1="296" x2="360" y2="344" stroke="#516274" stroke-width="5" marker-end="url(#arrow)" />
|
||||
|
||||
<rect x="48" y="360" width="624" height="142" rx="22" fill="#eef3f7" stroke="#738195" stroke-width="3" />
|
||||
<circle cx="100" cy="422" r="30" fill="#516274" />
|
||||
<text class="number" x="92" y="431">2</text>
|
||||
<text class="step" x="158" y="414">Outbox и публикация</text>
|
||||
<text class="copy" x="158" y="455">Есть ли message до worker?</text>
|
||||
<line x1="360" y1="504" x2="360" y2="552" stroke="#516274" stroke-width="5" marker-end="url(#arrow)" />
|
||||
|
||||
<rect x="48" y="568" width="624" height="154" rx="22" fill="#fff4d7" stroke="#c68819" stroke-width="3" />
|
||||
<circle cx="100" cy="636" r="30" fill="#b47700" />
|
||||
<text class="number" x="92" y="645">3</text>
|
||||
<text class="step" x="158" y="625">Журнал worker</text>
|
||||
<text class="copy" x="158" y="666">received → result → ack</text>
|
||||
<text class="small" x="158" y="697">Сравнить порядок и attempt</text>
|
||||
<line x1="360" y1="724" x2="360" y2="772" stroke="#516274" stroke-width="5" marker-end="url(#arrow)" />
|
||||
|
||||
<rect x="48" y="788" width="624" height="142" rx="22" fill="#e4f6ea" stroke="#25855a" stroke-width="3" />
|
||||
<circle cx="100" cy="850" r="30" fill="#16764c" />
|
||||
<text class="number" x="92" y="859">4</text>
|
||||
<text class="step" x="158" y="842">resultKey найден</text>
|
||||
<text class="copy" x="158" y="883">Повторный delivery только ACK</text>
|
||||
<line x1="360" y1="932" x2="360" y2="980" stroke="#516274" stroke-width="5" marker-end="url(#arrow)" />
|
||||
|
||||
<rect x="48" y="996" width="624" height="142" rx="22" fill="#fde9ea" stroke="#b4434b" stroke-width="3" />
|
||||
<circle cx="100" cy="1058" r="30" fill="#a6353d" />
|
||||
<text class="number" x="92" y="1067">5</text>
|
||||
<text class="step" x="158" y="1050">Постоянная ошибка</text>
|
||||
<text class="copy" x="158" y="1091">quarantined + DLX-разбор</text>
|
||||
</svg>
|
||||
|
After Width: | Height: | Size: 3.7 KiB |
@@ -0,0 +1,58 @@
|
||||
<svg xmlns="http://www.w3.org/2000/svg" viewBox="0 0 720 1320" role="img" aria-labelledby="title desc">
|
||||
<title id="title">Жизненный цикл фоновой задачи</title>
|
||||
<desc id="desc">Вертикальная схема показывает запись задачи и outbox, публикацию, обработку worker-ом, сохранение результата перед ack, повтор при временной ошибке и карантин постоянной ошибки.</desc>
|
||||
<defs>
|
||||
<marker id="arrow" viewBox="0 0 12 12" refX="6" refY="6" markerWidth="10" markerHeight="10" orient="auto-start-reverse"><path d="M 1 1 L 11 6 L 1 11 z" fill="#516274" /></marker>
|
||||
<style>
|
||||
.title { fill: #182230; font: 700 34px system-ui, -apple-system, sans-serif; }
|
||||
.lead { fill: #516274; font: 400 22px system-ui, -apple-system, sans-serif; }
|
||||
.step { fill: #182230; font: 700 27px system-ui, -apple-system, sans-serif; }
|
||||
.copy { fill: #354458; font: 400 23px system-ui, -apple-system, sans-serif; }
|
||||
.tag { fill: #ffffff; font: 700 22px system-ui, -apple-system, sans-serif; }
|
||||
.small { fill: #516274; font: 400 21px system-ui, -apple-system, sans-serif; }
|
||||
</style>
|
||||
</defs>
|
||||
<rect width="720" height="1320" fill="#f6f8fb" rx="28" />
|
||||
<text class="title" x="52" y="72">Жизненный цикл одной задачи</text>
|
||||
<text class="lead" x="52" y="108">jobId связывает запись, delivery и результат</text>
|
||||
|
||||
<rect x="48" y="148" width="624" height="142" rx="22" fill="#e7f0ff" stroke="#4a78bd" stroke-width="3" />
|
||||
<rect x="72" y="176" width="88" height="42" rx="21" fill="#3167b1" />
|
||||
<text class="tag" x="94" y="205">HTTP</text>
|
||||
<text class="step" x="184" y="203">Записать job и outbox</text>
|
||||
<text class="copy" x="184" y="244">Ответ: принято, jobId известен</text>
|
||||
<line x1="360" y1="292" x2="360" y2="340" stroke="#516274" stroke-width="5" marker-end="url(#arrow)" />
|
||||
|
||||
<rect x="48" y="356" width="624" height="142" rx="22" fill="#eef3f7" stroke="#738195" stroke-width="3" />
|
||||
<rect x="72" y="384" width="132" height="42" rx="21" fill="#516274" />
|
||||
<text class="tag" x="91" y="413">DISPATCH</text>
|
||||
<text class="step" x="226" y="411">Опубликовать jobId</text>
|
||||
<text class="copy" x="226" y="452">Только затем отметить outbox</text>
|
||||
<line x1="360" y1="500" x2="360" y2="548" stroke="#516274" stroke-width="5" marker-end="url(#arrow)" />
|
||||
|
||||
<rect x="48" y="564" width="624" height="142" rx="22" fill="#fff4d7" stroke="#c68819" stroke-width="3" />
|
||||
<rect x="72" y="592" width="112" height="42" rx="21" fill="#b47700" />
|
||||
<text class="tag" x="92" y="621">WORKER</text>
|
||||
<text class="step" x="210" y="619">Прочитать job по jobId</text>
|
||||
<text class="copy" x="210" y="660">Delivery ещё не подтверждён</text>
|
||||
<line x1="360" y1="708" x2="360" y2="756" stroke="#516274" stroke-width="5" marker-end="url(#arrow)" />
|
||||
|
||||
<rect x="48" y="772" width="624" height="156" rx="22" fill="#e4f6ea" stroke="#25855a" stroke-width="3" />
|
||||
<rect x="72" y="800" width="102" height="42" rx="21" fill="#16764c" />
|
||||
<text class="tag" x="97" y="829">УСПЕХ</text>
|
||||
<text class="step" x="200" y="827">Сохранить result + succeeded</text>
|
||||
<text class="copy" x="200" y="868">Только после этого отправить ACK</text>
|
||||
<text class="small" x="200" y="902">Повтор увидит terminal state</text>
|
||||
|
||||
<line x1="360" y1="930" x2="360" y2="978" stroke="#516274" stroke-width="5" marker-end="url(#arrow)" />
|
||||
<rect x="48" y="994" width="624" height="126" rx="22" fill="#fff1df" stroke="#c97918" stroke-width="3" />
|
||||
<rect x="72" y="1023" width="102" height="42" rx="21" fill="#b76500" />
|
||||
<text class="tag" x="95" y="1052">RETRY</text>
|
||||
<text class="step" x="201" y="1050">Временная ошибка: retry_wait</text>
|
||||
<text class="copy" x="201" y="1090">Ограниченный повтор вернётся к worker</text>
|
||||
|
||||
<line x1="360" y1="1122" x2="360" y2="1170" stroke="#516274" stroke-width="5" marker-end="url(#arrow)" />
|
||||
<rect x="48" y="1186" width="624" height="98" rx="22" fill="#fde9ea" stroke="#b4434b" stroke-width="3" />
|
||||
<text class="step" x="76" y="1228">Постоянная ошибка или лимит</text>
|
||||
<text class="copy" x="76" y="1265">→ quarantine / DLX</text>
|
||||
</svg>
|
||||
|
After Width: | Height: | Size: 4.4 KiB |
@@ -0,0 +1,55 @@
|
||||
<svg xmlns="http://www.w3.org/2000/svg" viewBox="0 0 720 1260" role="img" aria-labelledby="title desc">
|
||||
<title id="title">Граница повторов и подтверждения фоновой задачи</title>
|
||||
<desc id="desc">Вертикальная схема показывает, что worker сохраняет результат до подтверждения доставки, ограничивает временный повтор и отправляет постоянную ошибку в карантин без бесконечного requeue.</desc>
|
||||
<defs>
|
||||
<marker id="arrow" viewBox="0 0 12 12" refX="6" refY="6" markerWidth="10" markerHeight="10" orient="auto-start-reverse"><path d="M 1 1 L 11 6 L 1 11 z" fill="#516274" /></marker>
|
||||
<style>
|
||||
.title { fill: #182230; font: 700 34px system-ui, -apple-system, sans-serif; }
|
||||
.lead { fill: #516274; font: 400 22px system-ui, -apple-system, sans-serif; }
|
||||
.step { fill: #182230; font: 700 27px system-ui, -apple-system, sans-serif; }
|
||||
.copy { fill: #354458; font: 400 23px system-ui, -apple-system, sans-serif; }
|
||||
.tag { fill: #ffffff; font: 700 22px system-ui, -apple-system, sans-serif; }
|
||||
.small { fill: #516274; font: 400 21px system-ui, -apple-system, sans-serif; }
|
||||
</style>
|
||||
</defs>
|
||||
<rect width="720" height="1260" fill="#f6f8fb" rx="28" />
|
||||
<text class="title" x="50" y="72">Граница retry и ack</text>
|
||||
<text class="lead" x="50" y="108">delivery может повториться; jobId остаётся тем же</text>
|
||||
|
||||
<rect x="48" y="150" width="624" height="132" rx="22" fill="#e7f0ff" stroke="#4a78bd" stroke-width="3" />
|
||||
<rect x="72" y="178" width="112" height="42" rx="21" fill="#3167b1" />
|
||||
<text class="tag" x="92" y="207">DELIVERY</text>
|
||||
<text class="step" x="208" y="205">Получить tag и jobId</text>
|
||||
<text class="copy" x="208" y="246">Ack ещё не отправлен</text>
|
||||
<line x1="360" y1="284" x2="360" y2="332" stroke="#516274" stroke-width="5" marker-end="url(#arrow)" />
|
||||
|
||||
<rect x="48" y="348" width="624" height="146" rx="22" fill="#fff4d7" stroke="#c68819" stroke-width="3" />
|
||||
<rect x="72" y="376" width="116" height="42" rx="21" fill="#b47700" />
|
||||
<text class="tag" x="92" y="405">CHECK</text>
|
||||
<text class="step" x="212" y="403">Terminal state уже есть?</text>
|
||||
<text class="copy" x="212" y="444">Да: не делать effect второй раз</text>
|
||||
<text class="small" x="212" y="475">Нет: начать одну попытку</text>
|
||||
<line x1="360" y1="496" x2="360" y2="544" stroke="#516274" stroke-width="5" marker-end="url(#arrow)" />
|
||||
|
||||
<rect x="48" y="560" width="624" height="148" rx="22" fill="#e4f6ea" stroke="#25855a" stroke-width="3" />
|
||||
<rect x="72" y="588" width="104" height="42" rx="21" fill="#16764c" />
|
||||
<text class="tag" x="95" y="617">RESULT</text>
|
||||
<text class="step" x="201" y="615">resultKey + succeeded сохранены</text>
|
||||
<text class="copy" x="201" y="656">Только затем ACK этого tag</text>
|
||||
<text class="small" x="201" y="687">Порядок защищает от ранней потери</text>
|
||||
|
||||
<line x1="360" y1="710" x2="360" y2="758" stroke="#516274" stroke-width="5" marker-end="url(#arrow)" />
|
||||
<rect x="48" y="774" width="624" height="140" rx="22" fill="#fff1df" stroke="#c97918" stroke-width="3" />
|
||||
<rect x="72" y="802" width="98" height="42" rx="21" fill="#b76500" />
|
||||
<text class="tag" x="94" y="831">RETRY</text>
|
||||
<text class="step" x="196" y="829">Временная ошибка, attempt < 3</text>
|
||||
<text class="copy" x="196" y="870">retry_wait → задержанный повтор</text>
|
||||
<line x1="360" y1="916" x2="360" y2="964" stroke="#516274" stroke-width="5" marker-end="url(#arrow)" />
|
||||
|
||||
<rect x="48" y="980" width="624" height="154" rx="22" fill="#fde9ea" stroke="#b4434b" stroke-width="3" />
|
||||
<rect x="72" y="1008" width="156" height="42" rx="21" fill="#a6353d" />
|
||||
<text class="tag" x="110" y="1037">STOP</text>
|
||||
<text class="step" x="251" y="1035">Невалидно или лимит</text>
|
||||
<text class="copy" x="251" y="1076">quarantined → nack без requeue</text>
|
||||
<text class="small" x="251" y="1107">Дальше отдельный DLX-разбор</text>
|
||||
</svg>
|
||||
|
After Width: | Height: | Size: 4.2 KiB |
@@ -0,0 +1,415 @@
|
||||
function escapeHtml(value) {
|
||||
return String(value)
|
||||
.replaceAll('&', '&')
|
||||
.replaceAll('<', '<')
|
||||
.replaceAll('>', '>')
|
||||
.replaceAll('"', '"')
|
||||
.replaceAll("'", ''');
|
||||
}
|
||||
|
||||
function paragraph(text) {
|
||||
return '<p>' + text + '</p>';
|
||||
}
|
||||
|
||||
function heading(text) {
|
||||
return '<h2>' + text + '</h2>';
|
||||
}
|
||||
|
||||
function codeBlock(lines) {
|
||||
return '<pre><code>' + escapeHtml(lines.join('\n')) + '</code></pre>';
|
||||
}
|
||||
|
||||
function figure(src, alt, caption) {
|
||||
return '<figure><img src="' + src + '" alt="' + alt + '" loading="lazy" /><figcaption>' + caption + '</figcaption></figure>';
|
||||
}
|
||||
|
||||
function orderedList(items) {
|
||||
return '<ol>' + items.map((item) => '<li>' + item + '</li>').join('') + '</ol>';
|
||||
}
|
||||
|
||||
function dataTable(caption, headers, rows) {
|
||||
const head = '<thead><tr>' + headers.map((header) => '<th scope="col">' + header + '</th>').join('') + '</tr></thead>';
|
||||
const body = '<tbody>' + rows.map((row) => '<tr>' + row.map((cell) => '<td>' + cell + '</td>').join('') + '</tr>').join('') + '</tbody>';
|
||||
return '<div class="table-scroll"><table><caption>' + caption + '</caption>' + head + body + '</table></div>';
|
||||
}
|
||||
|
||||
function sourceList(items) {
|
||||
return '<ul>' + items.map((item) => '<li><a href="' + item.url + '" target="_blank" rel="noopener noreferrer">' + item.title + '</a> — ' + item.note + '</li>').join('') + '</ul>';
|
||||
}
|
||||
|
||||
function visibleText(html) {
|
||||
return html
|
||||
.replace(/<[^>]*>/g, ' ')
|
||||
.replaceAll(' ', ' ')
|
||||
.replaceAll('"', '"')
|
||||
.replaceAll(''', "'")
|
||||
.replaceAll('<', '<')
|
||||
.replaceAll('>', '>')
|
||||
.replaceAll('&', '&')
|
||||
.replace(/\s+/g, ' ')
|
||||
.trim();
|
||||
}
|
||||
|
||||
function proseText(html) {
|
||||
return visibleText(
|
||||
html
|
||||
.replace(/<pre><code>[\s\S]*?<\/code><\/pre>/g, '')
|
||||
.replace(/<figure>[\s\S]*?<\/figure>/g, '')
|
||||
.replace(/<div class="table-scroll">[\s\S]*?<\/div>/g, ''),
|
||||
);
|
||||
}
|
||||
|
||||
function createRevision(meta, bodyParts, sources) {
|
||||
const bodyHtml = bodyParts.join('\n');
|
||||
const proseLength = proseText(bodyHtml).length;
|
||||
|
||||
if (proseLength < 5000 || proseLength > 15000) {
|
||||
throw new Error(meta.slug + ': prose length must be 5000–15000, got ' + proseLength);
|
||||
}
|
||||
|
||||
if (sources.length < 2) {
|
||||
throw new Error(meta.slug + ': at least two primary or official sources are required');
|
||||
}
|
||||
|
||||
return {
|
||||
...meta,
|
||||
contentHtml: [bodyHtml, heading('Проверяемые источники'), sourceList(sources)].join('\n'),
|
||||
proseLength,
|
||||
};
|
||||
}
|
||||
|
||||
const amqpSpec = {
|
||||
title: 'AMQP 0-9-1 specification: basic.ack и basic.reject',
|
||||
url: 'https://www.rabbitmq.com/resources/specs/amqp-xml-doc0-9-1.pdf',
|
||||
note: 'первичная спецификация: delivery tag адресует доставку, а basic.reject с requeue управляет возвратом или отказом от сообщения',
|
||||
};
|
||||
|
||||
const rabbitAcknowledgements = {
|
||||
title: 'RabbitMQ: Consumer Acknowledgements and Publisher Confirms',
|
||||
url: 'https://www.rabbitmq.com/docs/3.13/confirms',
|
||||
note: 'официальное описание ручного ack, автоматического requeue не подтверждённой доставки и риска немедленного redelivery loop',
|
||||
};
|
||||
|
||||
const rabbitReliability = {
|
||||
title: 'RabbitMQ: Reliability guide',
|
||||
url: 'https://www.rabbitmq.com/docs/reliability',
|
||||
note: 'официальная граница между подтверждением доставки и обработкой, а также необходимость идемпотентного consumer при повторной доставке',
|
||||
};
|
||||
|
||||
const rabbitDlx = {
|
||||
title: 'RabbitMQ: Dead Letter Exchanges',
|
||||
url: 'https://www.rabbitmq.com/docs/next/dlx',
|
||||
note: 'официальное описание маршрутизации отклонённого сообщения в отдельный exchange и причин dead-lettering',
|
||||
};
|
||||
|
||||
const enqueueExample = [
|
||||
'// Учебный псевдокод: одна транзакция приложения, не клиент RabbitMQ.',
|
||||
'async function requestExport(input, db) {',
|
||||
' const jobId = makeStableId(input.accountId, input.period);',
|
||||
'',
|
||||
' await db.transaction(async (tx) => {',
|
||||
' await tx.insertJob({',
|
||||
" id: jobId, state: 'queued', attempt: 0, resultKey: null,",
|
||||
' });',
|
||||
' await tx.insertOutbox({',
|
||||
" type: 'report.export.requested', jobId: jobId, publishedAt: null,",
|
||||
' });',
|
||||
' });',
|
||||
'',
|
||||
' return { accepted: true, jobId: jobId };',
|
||||
'}',
|
||||
'',
|
||||
'// Отдельный dispatcher читает неопубликованный outbox.',
|
||||
'// Он отмечает publishedAt только после подтверждения выбранного broker client.',
|
||||
].join('\n');
|
||||
|
||||
const workerExample = [
|
||||
'// Учебный обработчик одного delivery. transport API намеренно абстрактный.',
|
||||
'async function handleDelivery(delivery, jobs, broker) {',
|
||||
' const job = await jobs.findForUpdate(delivery.jobId);',
|
||||
'',
|
||||
" if (job.state === 'succeeded' || job.state === 'quarantined') {",
|
||||
' await broker.ack(delivery.tag);',
|
||||
' return;',
|
||||
' }',
|
||||
'',
|
||||
' try {',
|
||||
" await jobs.markRunning(job.id, delivery.attempt);",
|
||||
' const resultKey = await writeReportOnce(job.id, job.payload);',
|
||||
' await jobs.markSucceeded(job.id, resultKey);',
|
||||
' await broker.ack(delivery.tag); // только после durable result + state',
|
||||
' } catch (error) {',
|
||||
' const nextAttempt = delivery.attempt + 1;',
|
||||
' await jobs.markRetryOrQuarantine(job.id, nextAttempt, error.code);',
|
||||
' await broker.nack(delivery.tag, { requeue: nextAttempt < 3 });',
|
||||
' }',
|
||||
'}',
|
||||
].join('\n');
|
||||
|
||||
const journalExample = [
|
||||
'{"event":"received","jobId":"export-42-2020-05","attempt":1,"redelivered":false}',
|
||||
'{"event":"retry_scheduled","jobId":"export-42-2020-05","attempt":1,"reason":"upstream_timeout"}',
|
||||
'{"event":"received","jobId":"export-42-2020-05","attempt":2,"redelivered":true}',
|
||||
'{"event":"result_saved","jobId":"export-42-2020-05","resultKey":"reports/export-42-2020-05.csv"}',
|
||||
'{"event":"ack_sent","jobId":"export-42-2020-05","attempt":2}',
|
||||
].join('\n');
|
||||
|
||||
const poisonJournalExample = [
|
||||
'{"event":"received","jobId":"export-43-2020-05","attempt":3,"redelivered":true}',
|
||||
'{"event":"validation_failed","jobId":"export-43-2020-05","reason":"unknown_report_kind"}',
|
||||
'{"event":"quarantined","jobId":"export-43-2020-05","queue":"jobs.quarantine"}',
|
||||
'{"event":"nack_sent","jobId":"export-43-2020-05","requeue":false}',
|
||||
].join('\n');
|
||||
|
||||
function runStateFixture() {
|
||||
const journal = [];
|
||||
const retryJob = { id: 'export-42-2020-05', state: 'queued', attempt: 0, resultKey: null };
|
||||
const poisonJob = { id: 'export-43-2020-05', state: 'queued', attempt: 0, resultKey: null };
|
||||
|
||||
function deliver(job, outcome) {
|
||||
job.attempt += 1;
|
||||
journal.push({ event: 'received', jobId: job.id, attempt: job.attempt, redelivered: job.attempt > 1 });
|
||||
|
||||
if (outcome === 'success') {
|
||||
job.resultKey = 'reports/' + job.id + '.csv';
|
||||
job.state = 'succeeded';
|
||||
journal.push({ event: 'result_saved', jobId: job.id, resultKey: job.resultKey });
|
||||
journal.push({ event: 'ack_sent', jobId: job.id, attempt: job.attempt });
|
||||
return;
|
||||
}
|
||||
|
||||
if (outcome === 'invalid' || job.attempt >= 3) {
|
||||
job.state = 'quarantined';
|
||||
journal.push({ event: 'quarantined', jobId: job.id, attempt: job.attempt, requeue: false });
|
||||
return;
|
||||
}
|
||||
|
||||
job.state = 'retry_wait';
|
||||
journal.push({ event: 'retry_scheduled', jobId: job.id, attempt: job.attempt, requeue: true });
|
||||
}
|
||||
|
||||
deliver(retryJob, 'transient');
|
||||
deliver(retryJob, 'success');
|
||||
deliver(poisonJob, 'transient');
|
||||
deliver(poisonJob, 'transient');
|
||||
deliver(poisonJob, 'invalid');
|
||||
|
||||
return {
|
||||
fixture: 'in-memory-job-state-machine',
|
||||
retryJob: retryJob,
|
||||
poisonJob: poisonJob,
|
||||
journal: journal,
|
||||
assertions: {
|
||||
ackAfterResult: retryJob.state === 'succeeded' && retryJob.resultKey !== null,
|
||||
poisonStopped: poisonJob.state === 'quarantined' && poisonJob.attempt === 3,
|
||||
retryWasRedelivered: journal.some((entry) => entry.jobId === retryJob.id && entry.redelivered === true),
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
const practiceArticle = createRevision(
|
||||
{
|
||||
slug: 'editorial-2020-05-practice-background-jobs',
|
||||
title: 'Фоновые задачи: сначала фиксируем намерение, потом запускаем worker',
|
||||
categories: ['Backend', 'Очереди', 'Практика'],
|
||||
cover: '/assets/editorial/2020/background-job-lifecycle-2020.svg',
|
||||
excerpt: 'Долгий экспорт не должен жить внутри HTTP-запроса. Собираем маленький контракт: запись задачи, публикация, обработка, сохранённый результат и ack в последнюю очередь.',
|
||||
readingMinutes: 14,
|
||||
},
|
||||
[
|
||||
paragraph('Симптом выглядит как «экспорт иногда исчезает». Пользователь нажал кнопку, HTTP-ответ вернул 202, но через десять минут нет ни файла, ни понятного статуса. Иногда всё хуже: worker успел записать отчёт, упал перед <code>ack</code>, а повтор создал второй файл или дважды отправил письмо. Цена не в самой минуте ожидания. Поддержка не может ответить, принята ли работа, разработчик не отличает потерю сообщения от дубля, а следующий срочный фикс добавляет ещё один таймаут вместо контракта.'),
|
||||
paragraph('В мае 2020 я бы не начинал с «фоновой платформы». Достаточна узкая связка delivery и backend: HTTP принимает намерение, приложение сохраняет состояние задачи, отдельный dispatcher публикует сообщение, worker делает работу, а очередь получает подтверждение только после сохранённого результата. Это учебная схема для одной операции <code>report.export</code>. Она не доказывает работу конкретного RabbitMQ-кластера и не обещает exactly-once; её задача — сделать каждую точку потери или повтора видимой.'),
|
||||
heading('У задачи есть свой владелец состояния'),
|
||||
paragraph('Сообщение в очереди не должно быть единственным местом, где живёт смысл работы. Брокер знает о delivery, но не обязан знать, создан ли CSV, обновлена ли строка в базе и что увидит пользователь. Поэтому у операции есть запись приложения с постоянным <code>jobId</code>, входными параметрами, числом попыток, текущим состоянием, ключом результата и короткой причиной последней неудачи. Один и тот же <code>jobId</code> проходит HTTP, outbox, message, worker и журнал. Это не трассировка всей системы, а минимальная нить для одного расследования.'),
|
||||
paragraph('Нормальный ответ HTTP — не «готово», а «принято»: <code>{ accepted: true, jobId }</code>. Клиент затем спрашивает статус именно этой записи или получает уведомление привычным для продукта способом. Если запись не создана, нет принятой задачи. Если она создана, но сообщение ещё не опубликовано, это отдельное наблюдаемое состояние, а не повод сказать пользователю, что worker уже начал работу.'),
|
||||
dataTable(
|
||||
'Контракт одной фоновой задачи: состояние не прячется внутри очереди',
|
||||
['Слой', 'Что хранит', 'Что считается доказательством', 'Чего не обещает'],
|
||||
[
|
||||
['HTTP', '<code>jobId</code> и ответ 202', 'в транзакции появилась запись задачи', 'что обработка уже завершена'],
|
||||
['База приложения', 'state, attempt, resultKey, lastError', 'задача читается по <code>jobId</code>', 'что сообщение дошло до broker'],
|
||||
['Outbox', 'событие публикации и признак отправки', 'есть элемент, который можно дочитать после рестарта', 'атомарный commit с внешним broker'],
|
||||
['Очередь', 'delivery выбранному consumer', 'ручной ack или nack по delivery tag', 'идемпотентность бизнес-операции'],
|
||||
['Worker', 'результат и переход состояния', 'resultKey сохранён до ack', 'безопасность внешнего побочного эффекта без ключа'],
|
||||
],
|
||||
),
|
||||
paragraph('Такой список убирает опасное сокращение «задача в очереди». У нас есть как минимум запись намерения, запись на публикацию, delivery и бизнес-результат. Они могут находиться в разных состояниях одновременно. Дежурный не должен гадать по отсутствию файла: он открывает строку задачи, смотрит <code>state</code>, <code>attempt</code> и <code>lastError</code>, после чего знает, на какой границе продолжать проверку.'),
|
||||
figure('/assets/editorial/2020/background-job-lifecycle-2020.svg', 'Вертикальная схема жизненного цикла фоновой задачи: HTTP сохраняет задачу и outbox, dispatcher публикует сообщение, worker сохраняет результат, затем отправляет ack; ошибочная задача уходит в повтор или карантин', 'Жизненный цикл показывает четыре разных факта: намерение записано, сообщение опубликовано, результат сохранён, delivery подтверждён. Ack стоит последним, потому что не заменяет бизнес-результат.'),
|
||||
heading('Разрыв между записью и публикацией надо назвать'),
|
||||
paragraph('Наивная последовательность «сначала записали заявку, потом отправили сообщение» ломается при падении между двумя строками. В базе уже есть <code>queued</code>, а delivery нет. Обратный порядок не лучше: worker может получить сообщение и не найти ещё незафиксированную запись. Распределённую транзакцию между базой и broker я здесь не предлагаю: для небольшого сервиса она быстро становится дороже самой задачи. Вместо этого полезнее в одной транзакции приложения сохранить и задачу, и outbox-запись.'),
|
||||
paragraph('После commit отдельный короткий dispatcher ищет outbox без <code>publishedAt</code>, передаёт минимальное сообщение <code>{ jobId, type, version }</code> выбранному клиенту broker и отмечает публикацию только по контракту этого клиента. Если процесс умер до такой отметки, dispatcher попробует снова. Отсюда следует неприятный, но здоровый вывод: consumer обязан выдержать duplicate. Повтор публикации — не дефект, который можно «выключить» флагом; это цена за возможность восстановиться после неясного обрыва.'),
|
||||
codeBlock(enqueueExample.split('\n')),
|
||||
paragraph('В примере нет SQL-схемы реального проекта и нет вызова библиотеки RabbitMQ. Важно другое: <code>jobId</code> создаётся до публикации, а outbox живёт рядом с бизнес-записью. В результате можно отдельно проверить, почему dispatcher не двигает запись: отсутствует ли соединение, неверен ли маршрут, не получено ли ожидаемое подтверждение или просто нет самого outbox-события. Это намного полезнее, чем повторно запускать весь экспорт по кнопке.'),
|
||||
heading('Ack подтверждает delivery, а не желание верить в успех'),
|
||||
paragraph('В AMQP delivery tag относится к конкретной доставке на конкретном канале, а не к вечному идентификатору задачи. Ручной <code>ack</code> говорит broker, что consumer принял ответственность за это delivery; после него broker может удалить сообщение. Поэтому ack до записи результата опасен: процесс может упасть после подтверждения, а нужная работа исчезнет из очереди. Ack после результата допускает обратное окно: результат есть, а ack не успел уйти. Тогда сообщение будет доставлено повторно. Это ожидаемый сценарий, не исключение.'),
|
||||
paragraph('Проверка порядка должна быть предельно приземлённой. До <code>ack</code> в хранилище уже есть <code>state = succeeded</code> и устойчивый <code>resultKey</code>. При повторной доставке worker читает эту запись, не запускает экспорт снова и подтверждает только повторное delivery. Если результат создаётся во внешнем сервисе, например объектном хранилище или email-шлюзе, одного флага <code>succeeded</code> недостаточно: нужен стабильный ключ объекта или ключ идемпотентности на стороне такого вызова. В этой статье мы ограничиваемся одной задачей и явно не выдаём этот принцип за универсальную гарантию всех сторонних систем.'),
|
||||
heading('Повтор — это отдельное состояние, не бесконечный requeue'),
|
||||
paragraph('У временной ошибки есть диагностируемая причина: короткий сбой зависимости, временный лимит или неготовый вход. Для неё можно сохранить <code>retry_wait</code>, номер попытки и код причины, а затем вернуть работу через заданную задержку и ограниченный маршрут. У невалидного входа причина другая: новый delivery не исправит неизвестный тип отчёта или отсутствующий обязательный параметр. Если каждое такое сообщение немедленно отправлять с <code>requeue: true</code>, worker будет тратить CPU на один и тот же отказ, а полезные задачи окажутся позади него.'),
|
||||
dataTable(
|
||||
'Минимальная retry-политика для учебного экспорта',
|
||||
['Наблюдение', 'Состояние задачи', 'Действие с delivery', 'Что остаётся для разбора'],
|
||||
[
|
||||
['Сохранён CSV и resultKey', '<code>succeeded</code>', '<code>ack</code>', 'jobId, ключ результата, попытка'],
|
||||
['Короткий timeout зависимости, попытка 1–2', '<code>retry_wait</code>', 'вернуть в задержанный retry-маршрут', 'код ошибки и следующая попытка'],
|
||||
['Невалидный payload или попытка 3', '<code>quarantined</code>', '<code>nack/reject</code> без requeue; DLX при настроенном маршруте', 'payload-версия, причина, jobId, delivery metadata'],
|
||||
['Повтор delivery после сбоя до ack', '<code>succeeded</code> уже есть', '<code>ack</code> без нового экспорта', 'признак redelivered и прежний resultKey'],
|
||||
],
|
||||
),
|
||||
paragraph('Poison message здесь не мистическая категория broker. Это сообщение, которое снова и снова не может пройти известный consumer-контракт. Его путь должен заканчиваться в отдельной карантинной очереди или другой управляемой поверхности, а не возвращаться в тот же hot loop. Карантин не означает «удалить и забыть»: для записи должны остаться <code>jobId</code>, версия payload, причина, число попыток и понятный владелец решения — исправить данные, исправить worker или осознанно отменить работу.'),
|
||||
heading('Маршрут проверки до первого реального запуска'),
|
||||
orderedList([
|
||||
'Записать для операции постоянный <code>jobId</code>, разрешённые состояния и условие, после которого пользователь может увидеть результат.',
|
||||
'В одной транзакции приложения создать job и outbox; после имитации рестарта проверить, что непросланный outbox всё ещё читается.',
|
||||
'На учебной фикстуре прогнать временную ошибку: первая попытка переходит в <code>retry_wait</code>, вторая сохраняет результат, и только затем появляется <code>ack_sent</code>.',
|
||||
'Отдельно прогнать невалидный payload или лимит попыток: запись становится <code>quarantined</code>, а маршруту не разрешён немедленный requeue.',
|
||||
'Повторить delivery для уже <code>succeeded</code> job и убедиться, что результат не создаётся второй раз, а worker подтверждает новое delivery.',
|
||||
'До подключения broker выбрать конкретный клиентский API, проверить его publisher-confirm и DLX-настройки на стенде; этот текст не подменяет такой прогон.',
|
||||
]),
|
||||
heading('Граница этой практики'),
|
||||
paragraph('Здесь нет настоящего очередного сервера, зарегистрированного consumer или production-лога. В модуле пакета есть только детерминированная in-memory фикстура состояний: она доказывает, что авторский переход «повтор → успех» отправляет ack после result и что третий неуспех переводит другую задачу в карантин. Она не доказывает поведение сети, задержку broker, порядок в нескольких worker и конфигурацию dead-letter exchange. Эти свойства должны быть проверены отдельным стендом с тем broker и клиентом, которые выбрал проект.'),
|
||||
],
|
||||
[amqpSpec, rabbitAcknowledgements, rabbitReliability, rabbitDlx],
|
||||
);
|
||||
|
||||
const mechanismArticle = createRevision(
|
||||
{
|
||||
slug: 'editorial-2020-05-mechanism-background-jobs',
|
||||
title: 'Под капотом фоновой задачи: delivery, ack, повтор и карантин',
|
||||
categories: ['Backend', 'Очереди', 'Механизмы'],
|
||||
cover: '/assets/editorial/2020/background-job-retry-boundary-2020.svg',
|
||||
excerpt: 'Разделяем бизнес-состояние задачи и одно delivery: где ставить ack, почему повтор нормален, как сделать одну задачу идемпотентной и остановить poison message.',
|
||||
readingMinutes: 15,
|
||||
},
|
||||
[
|
||||
paragraph('Симптом механической ошибки обычно звучит так: «мы же уже обработали сообщение, почему оно пришло ещё раз?» Или наоборот: worker вызвал API, после чего задача пропала без результата. Цена обоих случаев одна — команда смешала два разных объекта: бизнес-операцию и delivery от broker. Первый живёт столько, сколько нужен отчёт или запись; второе живёт до ack/nack на конкретном канале. Пока между ними нет явного перехода, дубль кажется аварией, а ранний ack кажется оптимизацией.'),
|
||||
paragraph('В мае 2020 полезнее освоить небольшой автомат, чем рисовать сложную оркестрацию. У операции есть <code>jobId</code>, у каждой доставки — свой tag и признак redelivery, а worker делает четыре действия в фиксированном порядке: получить delivery, захватить состояние задачи, сохранить исход работы, сообщить broker результат обработки. Этот порядок не даёт «ровно один раз» на всём мире. Он даёт проверяемую at-least-once модель, в которой повтор не обязан повторять эффект одной уже завершённой задачи.'),
|
||||
heading('Delivery и задача отвечают на разные вопросы'),
|
||||
paragraph('Delivery отвечает на вопрос broker: кто сейчас несёт ответственность за эту копию сообщения? Для ручного подтверждения ответ заканчивается <code>ack(tag)</code> или отрицательным ответом. Если соединение consumer закрылось до ack, broker может вернуть непроверенную доставку и позже отдать её тому же или другому worker. Задача отвечает на вопрос приложения: что пользователь попросил, на какой попытке это находится и где результат. Её нельзя определить только по тому, есть ли сообщение в очереди.'),
|
||||
paragraph('Из этой разницы получается рабочее правило. Не используем delivery tag как <code>jobId</code>: tag привязан к каналу и меняется при следующей доставке. В payload кладём короткий стабильный идентификатор и версию контракта, а детали входа храним у владельца задачи либо подписываем настолько ясно, чтобы worker мог проверить их версию. Когда приходит redelivery, worker ищет тот же <code>jobId</code>, а не пытается угадать, была ли такая строка по совпадению времени или имени файла.'),
|
||||
dataTable(
|
||||
'Два идентификатора и две ответственности',
|
||||
['Сущность', 'Живёт где', 'Меняется когда', 'Правильное применение'],
|
||||
[
|
||||
['<code>jobId</code>', 'в базе приложения, payload и журнале', 'не меняется между повторами', 'идемпотентность, статус, результат, поддержка'],
|
||||
['delivery tag', 'в канале broker/client', 'при каждом новом delivery', 'ровно один ack/nack конкретной доставки'],
|
||||
['attempt', 'в записи задачи или явно в повторном сообщении', 'при принятом решении о повторе', 'лимит retry и отчёт о причине'],
|
||||
['resultKey', 'в устойчивом storage/БД', 'один раз при успехе', 'доказательство, что effect уже готов до ack'],
|
||||
],
|
||||
),
|
||||
paragraph('Префикс «at least once» относится к доставке, не к тому, что пользователь увидит два отчёта. Повтор возможен и после publisher-side неопределённости, и после consumer-side падения. Поэтому нельзя надеяться, что флаг <code>redelivered</code> всё решит: он полезен для журнала и приоритета проверки, но не заменяет запись состояния. Безопасное решение на уровне одной задачи — выбирать один стабильный эффект: например, файл всегда пишется в <code>reports/{jobId}.csv</code>, а строка результата хранит этот же ключ.'),
|
||||
figure('/assets/editorial/2020/background-job-retry-boundary-2020.svg', 'Вертикальная схема границы повторов: worker получает delivery, сохраняет результат и только затем ack; временная ошибка идёт в ограниченный retry, а невалидная задача после лимита оказывается в карантине без бесконечного requeue', 'Повтор относится к delivery, а состояние — к jobId. На каждом пути сначала сохраняется решение приложения, затем broker получает ack или nack.'),
|
||||
heading('Ack ставим после устойчивого бизнес-результата'),
|
||||
paragraph('Ручной ack — это не запись в лог «worker начал работу». Это граница, после которой broker вправе удалить delivery. Если handler отправил ack сразу после получения, то сбой в SQL, файловом хранилище или внешнем API оставит ложный след: очередь считает задачу завершённой, приложение — нет. Если ack ставится после результата, возможен другой порядок: результат уже сохранён, а сеть закрылась. Следующее delivery должно обнаружить завершённую задачу и не делать effect второй раз.'),
|
||||
paragraph('Здесь важно назвать настоящий порядок устойчивости. Внутри приложения <code>markSucceeded</code> и <code>resultKey</code> должны быть сохранены одной согласованной операцией или в таком порядке, который можно восстановить после рестарта. Для внешнего вызова нужен собственный ключ: запись файла по детерминированному пути, уникальный ключ отправки или API-контракт, который принимает idempotency key. Нельзя сначала сделать необратимый эффект, а потом надеяться, что локальная таблица спасёт от дубля. Таблица лишь даёт worker решение при следующем delivery.'),
|
||||
codeBlock(workerExample.split('\n')),
|
||||
paragraph('Это не интерфейс конкретной Node-библиотеки. Он намеренно показывает точки, которые нельзя переставлять: terminal state проверяется до действия; попытка фиксируется до повторного маршрута; success пишется до ack; карантин фиксируется до <code>nack</code> без requeue. Реальный API может называть методы иначе и по-разному задавать задержку. До внедрения надо сверить эти места с версией выбранного client и с тем, настроен ли у очереди путь dead-lettering.'),
|
||||
heading('Повтор классифицируем до того, как вернуть сообщение'),
|
||||
paragraph('Не всякая ошибка заслуживает retry. Timeout до получения ответа иногда временный, но он не доказывает, что внешняя операция не состоялась; здесь особенно нужен ключ идемпотентности. Ошибка в payload, неизвестная версия сообщения или нарушенное обязательное поле повтором не исправится. Ошибка локальной валидации должна сразу закончить processing как карантин или осознанная отмена, а не навсегда держать одно delivery на голове очереди.'),
|
||||
dataTable(
|
||||
'Решение для ошибки одного delivery',
|
||||
['Класс', 'Пример симптома', 'Проверка перед решением', 'Следующее действие'],
|
||||
[
|
||||
['Успех', 'resultKey уже сохранён', 'состояние terminal и эффект читается по jobId', '<code>ack</code>; при дубле — только ack'],
|
||||
['Временная', 'короткий timeout до ответного байта', 'записаны attempt и причина; есть бюджет меньше лимита', 'перевести job в retry_wait, вернуть через ограниченный маршрут'],
|
||||
['Неопределённая внешняя', 'таймаут после отправки запроса', 'проверить внешний ключ или статус по jobId', 'не создавать второй эффект; повторять только через идемпотентный контракт'],
|
||||
['Постоянная', 'payload не проходит version/validation', 'причина воспроизводится на тех же данных', 'quarantine и nack/reject без requeue'],
|
||||
],
|
||||
),
|
||||
paragraph('Задержка между попытками тоже часть контракта. В 2020 году для одного сервиса можно обойтись простой retry-очередью с TTL или расписанием, которое уже умеет конкретная библиотека; не обязательно строить платформу. Но у каждой задачи должны быть фиксированные максимум попыток, причина последнего перехода и время следующего допуска. Формула exponential backoff не спасает, если неизвестно, какая именно попытка уже была и кто вернёт сообщение из задержки.'),
|
||||
heading('Poison message — это работа, которая больше не должна горячо крутиться'),
|
||||
paragraph('В AMQP <code>basic.reject</code> с <code>requeue=false</code> не обозначает «успех». Оно заканчивает обработку данного delivery; при заранее настроенном dead-letter exchange broker направляет сообщение на отдельный маршрут, иначе оно может быть отброшено. Это требует явной проектной договорённости: куда попадёт карантин, кто посмотрит его, как сопоставить payload с записью задачи и как защитить чувствительные поля от попадания в журнал.'),
|
||||
paragraph('Не стоит строить логику на бесконечном <code>nack(requeue=true)</code>. RabbitMQ прямо предупреждает, что consumer, который все время возвращает delivery, может создать затратный redelivery loop. Лимит попыток в записи задачи плюс отдельный результат <code>quarantined</code> разрывают этот цикл с понятной ценой: часть работы не завершена автоматически, зато полезные сообщения продолжают получать worker, а человек получает конкретную причину вместо бесконечного шума.'),
|
||||
heading('Короткий журнал заменяет догадку о порядке'),
|
||||
paragraph('Для первой версии не нужна отдельная observability-платформа. Достаточно, чтобы каждый переход писал один и тот же набор: <code>event</code>, <code>jobId</code>, <code>attempt</code>, признак <code>redelivered</code>, причина и при успехе <code>resultKey</code>. Тогда вопрос «был ли ack до результата?» проверяется порядком пяти строк, а не памятью того, кто разбирает сбой. Нельзя класть в такой журнал полный payload, токены или документ пользователя: для связи достаточно идентификатора и безопасной классификации ошибки.'),
|
||||
codeBlock(journalExample.split('\n')),
|
||||
paragraph('Это пример ожидаемого учебного журнала, не снятый log реального worker. Первая строка показывает начало первой попытки, вторая — сохранённое решение о повторе. Вторая доставка имеет тот же <code>jobId</code>, но уже другой delivery и <code>redelivered: true</code>. Только после строки <code>result_saved</code> возникает <code>ack_sent</code>. Если эти две строки поменялись местами, расследование закончено: worker подтверждает работу до того, как может доказать её эффект.'),
|
||||
heading('Маршрут проверки автомата'),
|
||||
orderedList([
|
||||
'Составить таблицу состояний для одного <code>jobId</code>: queued, running, retry_wait, succeeded и quarantined; убрать неявное «вроде выполняется».',
|
||||
'Взять один payload и показать два delivery с разными tags; убедиться, что оба ищут одну и ту же запись задачи.',
|
||||
'Смоделировать падение после <code>result_saved</code> до ack. Повтор должен только подтвердить новое delivery, а не записать второй результат.',
|
||||
'Смоделировать временную ошибку до лимита и проверить, что attempt и причина сохранены до постановки retry.',
|
||||
'Смоделировать невалидный payload либо третий отказ: задача становится quarantined, а немедленный requeue не разрешён.',
|
||||
'На отдельном стенде выбранного broker проверить реальную семантику ack/nack, redelivery, policy DLX и задержки; учебный автомат этого не заменяет.',
|
||||
]),
|
||||
heading('Граница модели'),
|
||||
paragraph('Этот материал не утверждает, что любая очередь предоставляет одинаковый delivery count, delayed retry или безопасный dead-letter путь. В нём использованы понятия AMQP 0-9-1 и документация RabbitMQ как проверяемый ориентир; конкретная конфигурация зависит от версии broker, типа очереди и client library. Внутри пакета выполнена только детерминированная фикстура переходов в памяти. Она не открывает соединение с RabbitMQ, не измеряет throughput, не запускает конкурирующих worker и не заменяет интеграционный тест проекта.'),
|
||||
],
|
||||
[amqpSpec, rabbitAcknowledgements, rabbitReliability, rabbitDlx],
|
||||
);
|
||||
|
||||
const fieldArticle = createRevision(
|
||||
{
|
||||
slug: 'editorial-2020-05-field-background-jobs',
|
||||
title: 'Разбор: почему экспорт отчёта повторился и как остановить poison message',
|
||||
categories: ['Backend', 'Очереди', 'Разбор'],
|
||||
cover: '/assets/editorial/2020/background-job-diagnosis-2020.svg',
|
||||
excerpt: 'Учебный разбор двух путей: результат создан до потерянного ack и невалидная задача крутится в requeue. Собираем журнал, меняем порядок и вводим карантин.',
|
||||
readingMinutes: 14,
|
||||
},
|
||||
[
|
||||
paragraph('Симптом в учебном разборе такой: пользователь запрашивает экспорт, а через несколько минут видит два одинаковых файла. Одновременно другая заявка с неизвестным типом отчёта снова и снова появляется у worker, занимая очередь. Цена двойная. Первый сбой создаёт лишний внешний эффект и спор, какая копия верная; второй забирает время worker и скрывает полезные задачи под повторяющейся ошибкой. Фраза «очередь доставила дважды» описывает факт, но ещё не называет место, где принято неверное решение.'),
|
||||
paragraph('Разберу не production-инцидент, а анонимизированную in-memory фикстуру мая 2020 года. В ней нет реального broker, файлового хранилища, user data или измеренной нагрузки. Зато есть две детерминированные цепочки с одним <code>jobId</code> каждая: временный отказ между сохранением результата и ack, а также невалидный payload после лимита попыток. Цель — показать практическое расследование: симптом → причина → проверка → действие, а не рассказать историю успеха постфактум.'),
|
||||
heading('Сначала отделяем факт результата от факта delivery'),
|
||||
paragraph('Первый экспорт имеет <code>jobId = export-42-2020-05</code>. Worker получил delivery, записал CSV по устойчивому ключу и пометил задачу как <code>succeeded</code>. Затем соединение до broker оборвалось до ack. У broker остаётся непроверенное delivery, поэтому следующий worker получает ту же бизнес-задачу повторно. Если handler относится к любому received message как к новому, он снова вызывает экспорт и пишет второй файл с новым случайным именем. Это не исправляется большим timeout: проблема в том, что idempotency check находится после эффекта или отсутствует.'),
|
||||
paragraph('Вторая заявка — <code>export-43-2020-05</code> — содержит неизвестный <code>reportKind</code>. Worker ловит ошибку, делает <code>nack(requeue=true)</code> и тут же получает ту же доставку снова. Никакая пауза не сделает неизвестный тип валидным. Пока задача не имеет состояния <code>quarantined</code> и ограничителя попыток, очередь по сути работает как генератор одинаковых ошибок. Здесь цена уже операционная: журнал растёт, полезная работа ждёт, а владелец данных не получает короткий список того, что нужно исправить.'),
|
||||
codeBlock(poisonJournalExample.split('\n')),
|
||||
paragraph('Строки выше — синтетический журнал фикстуры, не вывод запущенного RabbitMQ consumer. Они важны именно порядком. У poison-задачи третья попытка ещё фиксирует вход и причину, затем приложение сохраняет <code>quarantined</code>, и только после этого выбранному transport посылается отрицательный ответ без requeue. В реальном AMQP дальнейшая судьба зависит от настроенного dead-letter exchange: без маршрута сообщение может быть отброшено. Поэтому «карантин» обязан существовать не только как слово в коде, но и как проверяемая конфигурация выбранного окружения.'),
|
||||
dataTable(
|
||||
'Карта расследования: какой факт исключает какую гипотезу',
|
||||
['Наблюдение', 'Причина, которую проверяем', 'Минимальное доказательство', 'Действие'],
|
||||
[
|
||||
['Есть resultKey, но нет ack_sent', 'сбой произошёл в узком окне после результата', 'job.state = succeeded раньше следующего received', 'при повторе не создавать результат, подтвердить новое delivery'],
|
||||
['Два файла для одного jobId', 'внешний эффект не имеет стабильного ключа либо check сделан поздно', 'сопоставить имена файлов и порядок journal', 'путь результата построить из jobId, terminal state читать до export'],
|
||||
['attempt растёт, причина одна и та же', 'постоянный payload повторно requeue', 'validation_failed повторяется на равных входных данных', 'пометить quarantined и направить в DLX/разбор'],
|
||||
['queued долго без received', 'outbox не опубликован либо worker не читает маршрут', 'есть job/outbox, но нет publish marker и journal delivery', 'проверить dispatcher, binding и конкретный broker client'],
|
||||
],
|
||||
),
|
||||
paragraph('Эта таблица не заменяет доступ к очереди. Она задаёт порядок вопросов до изменения кода. Если уже есть <code>resultKey</code>, не нужно стартовать новый экспорт «на всякий случай». Если причина <code>unknown_report_kind</code> воспроизводится из сохранённой версии payload, не нужно увеличивать retry. Если у задачи нет <code>received</code>, бесполезно рассматривать handler: сперва ищем outbox, публикацию и маршрут. Каждый шаг привязывает действие к одному наблюдаемому факту.'),
|
||||
figure('/assets/editorial/2020/background-job-diagnosis-2020.svg', 'Вертикальная схема диагностики фоновой задачи: по jobId проверяют запись задачи, outbox и сообщение, затем журнал worker, сохранённый resultKey и при постоянной ошибке карантинную очередь', 'Разбор идёт не по названию компонента, а по пути одного jobId. Это уменьшает риск перезапустить эффект, когда проблема находится до worker или после результата.'),
|
||||
heading('Причина первого дубля: случайное имя и ack не на той стороне'),
|
||||
paragraph('Плохой вариант handler выглядит почти естественно: он берёт сообщение, сразу начинает генерацию, формирует имя из текущего времени, отправляет ack и только затем пытается отметить успех. В нём две точки неопределённости. Во-первых, повтору нечем доказать, что прежняя генерация уже завершилась: название файла другое, а состояние ещё не terminal. Во-вторых, ack может добраться до broker раньше записи статуса. При сбое получаем либо потерянную работу, либо повтор без защиты.'),
|
||||
paragraph('Исправление на уровне одной задачи не требует общего дедупликатора. Worker сначала читает строку по <code>jobId</code> с блокировкой, проверяет terminal states и резервирует попытку. Результат записывается по детерминированному ключу <code>reports/{jobId}.csv</code>. После записи в этом же бизнес-шаге сохраняются <code>resultKey</code> и <code>succeeded</code>. Если после этого broker повторно доставит сообщение, handler видит terminal state, не пишет файл заново и только завершает текущий delivery. Для email или внешнего API нужно отдельно убедиться, что принимающая сторона поддерживает такой ключ; путь файла не решает чужой side effect.'),
|
||||
dataTable(
|
||||
'Изменение порядка для повторной доставки',
|
||||
['Старая последовательность', 'Риск', 'Новая последовательность', 'Проверяемый результат'],
|
||||
[
|
||||
['receive → generate random file → ack → save state', 'дубль или потеря при падении между шагами', 'receive → read job → save deterministic result + succeeded → ack', 'повтор видит succeeded и не создаёт второй файл'],
|
||||
['catch → nack(requeue=true) всегда', 'горячий цикл на невалидном payload', 'classify → retry_wait или quarantined → nack по решению', 'attempt ограничен, причина остаётся рядом с jobId'],
|
||||
['искать ошибку по времени', 'непонятно, к какой попытке относится строка', 'писать jobId, attempt, event, reason', 'одна цепочка читается без догадки о совпадении'],
|
||||
],
|
||||
),
|
||||
paragraph('Здесь нет обещания, что SQL-блокировка сделает worker глобально одиночным. Она лишь защищает запись задачи в границе выбранной базы. Конкурирующие worker, внешнее хранилище и сеть всё равно требуют проверяемого контракта. Поэтому практический критерий короче: каждый новый delivery для уже <code>succeeded</code> обязан завершиться без нового результата. Если это нельзя проверить, слово «идемпотентность» в код-ревью пока ничего не означает.'),
|
||||
heading('Причина hot loop: постоянную ошибку приняли за временную'),
|
||||
paragraph('Повтор нужен, когда новое время может изменить исход: зависимость была недоступна, лимит снят, ожидаемая запись ещё не появилась. Но <code>unknown_report_kind</code> не зависит от времени. Для такого случая worker должен назвать ошибку постоянной, сохранить её в записи и завершить автоматический путь. В AMQP отрицательный ответ без requeue может направить сообщение в DLX, если проект это настроил. Отдельный маршрут делает ошибку предметом разбора, а не бесконечным consumer workload.'),
|
||||
paragraph('Карантин не стоит использовать как корзину для всех исключений. Сначала сохраняем тип причины и версию payload: это отделяет неисправимые данные от дефекта worker после обновления. Затем владелец решает: поправить данные и переиздать новую задачу с новым или тем же бизнес-ключом, починить consumer и вручную вернуть сообщение через контролируемый маршрут, либо отменить операцию. Автоматическое чтение карантина обратно в основную очередь без исправления причины снова создаёт тот же loop, только с более длинным названием.'),
|
||||
heading('Фикстура как маленький регрессионный контракт'),
|
||||
paragraph('В модуле ревизий есть команда <code>node scripts/upgrade-2020-05.mjs --verify-fixture</code>. Она не эмулирует AMQP frames. Она детерминированно создаёт две записи в памяти: первая переживает transient ошибку, получает повторную доставку, сохраняет result и ack; вторая после третьего неуспеха переходит в <code>quarantined</code>. Вывод — JSON-журнал и три булевых условия. Такой тест полезен тем, что не позволяет незаметно переставить <code>result_saved</code> и <code>ack_sent</code> в учебном алгоритме.'),
|
||||
codeBlock(journalExample.split('\n')),
|
||||
paragraph('Фикстура не даёт ложной уверенности в broker. У неё нет TCP-соединения, реального delivery tag, политики DLX, нескольких consumer или диска. Но она отделяет две логические проверки, которые можно выполнить без инфраструктуры: для retry есть новая попытка с тем же jobId, а успешный путь пишет result до ack; poison-путь обрывает requeue на известном пределе. После выбора библиотеки эту же пару сценариев нужно повторить на интеграционном стенде и сравнить реальные журналы с ожидаемыми переходами.'),
|
||||
heading('Маршрут разбора перед исправлением'),
|
||||
orderedList([
|
||||
'Взять один конкретный jobId и собрать рядом запись задачи, outbox, журнал worker, ключ результата и информацию о current attempt.',
|
||||
'Проверить, в каком порядке появились result_saved, succeeded и ack_sent; не делать новый экспорт, пока это не ясно.',
|
||||
'Для повтора сравнить jobId, а не delivery tag: новый tag не означает новую бизнес-операцию.',
|
||||
'Классифицировать последнюю ошибку как временную, неопределённую внешнюю или постоянную; записать основание рядом с attempt.',
|
||||
'Для постоянной ошибки остановить requeue, перевести job в quarantined и проверить, что выбранная DLX/карантинная поверхность действительно принимает сообщение.',
|
||||
'После изменения прогнать in-memory фикстуру, затем отдельный broker-интеграционный сценарий с падением до ack; в этом пакете выполнен только первый шаг.',
|
||||
]),
|
||||
heading('Граница полевого разбора'),
|
||||
paragraph('Все идентификаторы, причины и строки журнала здесь придуманы для проверки переходов. Нет реального файла, заказчика, очереди, RabbitMQ policy, production-config или browser-действия. Тексты опираются на спецификацию AMQP и официальную документацию RabbitMQ, чтобы не выдумывать смысл ack, reject и redelivery, но не выдают современную документацию за снимок конкретной инфраструктуры мая 2020 года. Автор этого периода умеет провести узкое backend/delivery расследование и оставить route для стенда; он ещё не заявляет готовую платформу наблюдаемости или сложную оркестрацию.'),
|
||||
],
|
||||
[amqpSpec, rabbitAcknowledgements, rabbitReliability, rabbitDlx],
|
||||
);
|
||||
|
||||
export const revisions = [practiceArticle, mechanismArticle, fieldArticle];
|
||||
|
||||
if (process.argv.includes('--print-revisions')) {
|
||||
process.stdout.write(JSON.stringify(revisions));
|
||||
} else if (process.argv.includes('--verify-fixture')) {
|
||||
const fixture = runStateFixture();
|
||||
if (!Object.values(fixture.assertions).every(Boolean)) {
|
||||
throw new Error('State fixture assertions failed');
|
||||
}
|
||||
process.stdout.write(JSON.stringify(fixture, null, 2) + '\n');
|
||||
}
|
||||
Reference in New Issue
Block a user