1"use strict";(self.webpackChunkrabbitmq_website=self.webpackChunkrabbitmq_website||[]).push([["12082"],{18231(e,n,s){s.r(n),s.d(n,{metadata:()=>i,default:()=>u,frontMatter:()=>c,contentTitle:()=>t,toc:()=>l,assets:()=>o});var i=JSON.parse('{"id":"consumer-prefetch","title":"Consumer Prefetch","description":"\x3c!--","source":"@site/versioned_docs/version-4.1/consumer-prefetch.md","sourceDirName":".","slug":"/consumer-prefetch","permalink":"/docs/4.1/consumer-prefetch","draft":false,"unlisted":false,"editUrl":"https://github.com/rabbitmq/rabbitmq-website/tree/main/versioned_docs/version-4.1/consumer-prefetch.md","tags":[],"version":"4.1","frontMatter":{"title":"Consumer Prefetch"},"sidebar":"docsSidebar","previous":{"title":"Consumer Cancellation Notifications","permalink":"/docs/4.1/consumer-cancel"},"next":{"title":"Consumer Priorities","permalink":"/docs/4.1/consumer-priority"}}'),r=s(74848),a=s(28453);let c={title:"Consumer Prefetch"},t="Consumer Prefetch",o={},l=[{value:"Overview",id:"overview",level:2},{value:"Single Consumer",id:"single-consumer",level:2},{value:"Independent Consumers",id:"independent-consumers",level:2},{value:"Multiple Consumers Sharing the Limit",id:"sharing-the-limit",level:2},{value:"Configurable Default Prefetch",id:"default-limit",level:2}];function h(e){let n={a:"a",code:"code",h1:"h1",h2:"h2",header:"header",p:"p",pre:"pre",...(0,a.R)(),...e.components};return(0,r.jsxs)(r.Fragment,{children:[(0,r.jsx)(n.header,{children:(0,r.jsx)(n.h1,{id:"consumer-prefetch",children:"Consumer Prefetch"})}),"\n",(0,r.jsx)(n.h2,{id:"overview",children:"Overview"}),"\n",(0,r.jsxs)(n.p,{children:["Consumer prefetch is an extension to the ",(0,r.jsx)(n.a,{href:"./confirms",children:"channel prefetch mechanism"}),"."]}),"\n",(0,r.jsxs)(n.p,{children:["AMQP 0-9-1 specifies the ",(0,r.jsx)(n.code,{children:"basic.qos"})," method to make it possible to\n",(0,r.jsx)(n.a,{href:"./confirms",children:"limit the number of unacknowledged messages"}),' on a channel (or\nconnection) when consuming (aka "prefetch count"). Unfortunately\nthe channel is not the ideal scope for this - since a single\nchannel may consume from multiple queues, the channel and the\nqueue(s) need to coordinate with each other for every message\nsent to ensure they don\'t go over the limit. This is slow on a\nsingle machine, and very slow when consuming across a cluster.']}),"\n",(0,r.jsx)(n.p,{children:"Furthermore for many uses it is simply more natural to specify\na prefetch count that applies to each consumer."}),"\n",(0,r.jsx)(n.p,{children:"Therefore RabbitMQ slightly deviates from the AMQP 0-9-1 spec\nwhen it comes to how the prefetch is applied to multiple consumers\non a channel:"}),"\n",(0,r.jsxs)("table",{class:"styled-table",children:[(0,r.jsxs)("tr",{children:[(0,r.jsxs)("th",{children:["Meaning of ",(0,r.jsx)("code",{children:"prefetch_count"})," in AMQP 0-9-1"]}),(0,r.jsxs)("th",{children:["Meaning of ",(0,r.jsx)("code",{children:"prefetch_count"})," in RabbitMQ"]})]}),(0,r.jsxs)("tr",{children:[(0,r.jsx)("td",{children:"shared across all consumers on the channel"}),(0,r.jsx)("td",{children:"applied separately to each new consumer on the channel"})]})]}),"\n",(0,r.jsx)(n.h2,{id:"single-consumer",children:"Single Consumer"}),"\n",(0,r.jsx)(n.p,{children:"The following basic example in Java will receive a maximum of 10\nunacknowledged messages at once:"}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-java",children:'Channel channel = ...;\nConsumer consumer = ...;\nchannel.basicQos(10); // Per consumer limit\nchannel.basicConsume("my-queue", false, consumer);\n'})}),"\n",(0,r.jsxs)(n.p,{children:["A value of ",(0,r.jsx)(n.code,{children:"0"})," is treated as infinite, allowing any number of unacknowledged\nmessages."]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-java",children:'Channel channel = ...;\nConsumer consumer = ...;\nchannel.basicQos(0); // No limit for this consumer\nchannel.basicConsume("my-queue", false, consumer);\n'})}),"\n",(0,r.jsx)(n.h2,{id:"independent-consumers",children:"Independent Consumers"}),"\n",(0,r.jsx)(n.p,{children:"This example starts two consumers on the same channel, each of\nwhich will independently receive a maximum of 10 unacknowledged\nmessages at once:"}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-java",children:'Channel channel = ...;\nConsumer consumer1 = ...;\nConsumer consumer2 = ...;\nchannel.basicQos(10); // Per consumer limit\nchannel.basicConsume("my-queue1", false, consumer1);\nchannel.basicConsume("my-queue2", false, consumer2);\n'})}),"\n",(0,r.jsx)(n.h2,{id:"sharing-the-limit",children:"Multiple Consumers Sharing the Limit"}),"\n",(0,r.jsxs)(n.p,{children:["The AMQP 0-9-1 specification does not explain what happens if you\ninvoke ",(0,r.jsx)(n.code,{children:"basic.qos"})," multiple times with different\n",(0,r.jsx)(n.code,{children:"global"})," values. RabbitMQ interprets this as meaning\nthat the two prefetch limits should be enforced independently of\neach other; consumers will only receive new messages when neither\nlimit on unacknowledged messages has been reached."]}),"\n",(0,r.jsx)(n.p,{children:"For example:"}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-java",children:'Channel channel = ...;\nConsumer consumer1 = ...;\nConsumer consumer2 = ...;\nchannel.basicQos(10, false); // Per consumer limit\nchannel.basicQos(15, true); // Per channel limit\nchannel.basicConsume("my-queue1", false, consumer1);\nchannel.basicConsume("my-queue2", false, consumer2);\n'})}),"\n",(0,r.jsx)(n.p,{children:"These two consumers will only ever have 15 unacknowledged\nmessages between them, with a maximum of 10 messages for each\nconsumer. This will be slower than the above examples, due to\nthe additional overhead of coordinating between the channel and\nthe queues to enforce the global limit."}),"\n",(0,r.jsx)(n.h2,{id:"default-limit",children:"Configurable Default Prefetch"}),"\n",(0,r.jsxs)(n.p,{children:["RabbitMQ can use a default prefetch that will be applied if the consumer doesn't specify one.\nThe value can be configured as ",(0,r.jsx)(n.code,{children:"rabbit.default_consumer_prefetch"})," in the ",(0,r.jsx)(n.a,{href:"./configure#advanced-config-file",children:"advanced configuration file"}),":"]}),"\n",(0,r.jsx)(n.pre,{children:(0,r.jsx)(n.code,{className:"language-erlang",children:"%% advanced.config file\n[\n {rabbit, [\n {default_consumer_prefetch, {false,250}}\n ]\n }\n].\n"})})]})}function u(e={}){let{wrapper:n}={...(0,a.R)(),...e.components};return n?(0,r.jsx)(n,{...e,children:(0,r.jsx)(h,{...e})}):h(e)}},28453(e,n,s){s.d(n,{R:()=>c,x:()=>t});var i=s(96540);let r={},a=i.createContext(r);function c(e){let n=i.useContext(a);return i.useMemo(function(){return"function"==typeof e?e(n):{...n,...e}},[n,e])}function t(e){let n;return n=e.disableParentContext?"function"==typeof e.components?e.components(r):e.components||r:c(e.components),i.createElement(a.Provider,{value:n},e.children)}}}]);
Line numbers count LF bytes from the start of the resource, as the search results do. Vendor segments are library code the classifier recognised; they are stored but not indexed. Bytes are shown as Latin1 characters, one per byte.