1"use strict";(self.webpackChunkmia_platform_docs=self.webpackChunkmia_platform_docs||[]).push([["14098"],{567440(e,n,t){t.r(n),t.d(n,{metadata:()=>s,default:()=>p,frontMatter:()=>i,contentTitle:()=>r,toc:()=>d,assets:()=>c});var s=JSON.parse('{"id":"products/fast_data_v2/mongezium_cdc/usage","title":"Usage","description":"Startup Behavior","source":"@site/versioned_docs/version-15.1.0/products/fast_data_v2/mongezium_cdc/30_Usage.md","sourceDirName":"products/fast_data_v2/mongezium_cdc","slug":"/products/fast_data_v2/mongezium_cdc/usage","permalink":"/docs/products/fast_data_v2/mongezium_cdc/usage","draft":false,"unlisted":false,"tags":[],"version":"15.1.0","sidebarPosition":30,"frontMatter":{"id":"usage","title":"Usage","sidebar_label":"Usage"},"sidebar":"fastDataV2","previous":{"title":"Configuration","permalink":"/docs/products/fast_data_v2/mongezium_cdc/configuration"},"next":{"title":"Metrics","permalink":"/docs/products/fast_data_v2/mongezium_cdc/metrics"}}'),o=t(474848),a=t(28453);let i={id:"usage",title:"Usage",sidebar_label:"Usage"},r,c={},d=[{value:"Startup Behavior",id:"startup-behavior",level:2},{value:"Usage Requirements",id:"usage-requirements",level:2},{value:"Messages Spec",id:"messages-spec",level:2}];function l(e){let n={a:"a",admonition:"admonition",code:"code",h2:"h2",img:"img",li:"li",p:"p",pre:"pre",ul:"ul",...(0,a.R)(),...e.components};return(0,o.jsxs)(o.Fragment,{children:[(0,o.jsx)(n.h2,{id:"startup-behavior",children:"Startup Behavior"}),"\n",(0,o.jsxs)(n.p,{children:["In the diagram below is described how the service decides whether to resume\na change stream or start a new one. Such behavior can be configured for\neach collection selecting the desired ",(0,o.jsx)(n.code,{children:"snapshot"})," mode."]}),"\n",(0,o.jsx)(n.p,{children:(0,o.jsx)(n.img,{alt:"mongezium_init_flow",src:t(201104).A+"",width:"4168",height:"1786"})}),"\n",(0,o.jsxs)(n.admonition,{type:"warning",children:[(0,o.jsxs)(n.p,{children:["To let Mongezium resume change stream events, the header of each message produced by Mongezium\ncontains a ",(0,o.jsx)(n.a,{href:"https://www.mongodb.com/docs/manual/changeStreams/#resume-tokens",children:"resume token"}),".\nAt startup, the service will read the latest message and create the change stream with the given resume token."]}),(0,o.jsxs)(n.p,{children:["It may happen that the ",(0,o.jsx)(n.code,{children:"oplog.rs"})," collection grows in size and certain resume token may disappear.\nIn such case, a new change stream will be opened, listening for new changes\nwithout performing any snapshot operation."]}),(0,o.jsxs)(n.p,{children:["To enforce the execution of a snapshot\nprocedure, please set the ",(0,o.jsx)(n.code,{children:"snapshot"})," field of collections entries to ",(0,o.jsx)(n.code,{children:"when_needed"}),".\nConsequently, when a resume token will return a not found error, a new snapshot\nwill be performed."]})]}),"\n",(0,o.jsx)(n.h2,{id:"usage-requirements",children:"Usage Requirements"}),"\n",(0,o.jsx)(n.p,{children:"To use the application, the following requirements must be met:"}),"\n",(0,o.jsxs)(n.ul,{children:["\n",(0,o.jsx)(n.li,{children:"MongoDB must be in replica-set."}),"\n",(0,o.jsxs)(n.li,{children:["the connection string must have privileges to access the ",(0,o.jsx)(n.code,{children:"oplog"})," and the ",(0,o.jsx)(n.code,{children:"admin"})," collection. More specifically, it needs permission to enable ",(0,o.jsx)(n.code,{children:"changeStreamPreAndPostImages"})," on the collection of the configured database;"]}),"\n",(0,o.jsxs)(n.li,{children:["Kafka connection must have permission to read/write the topics declared in the ",(0,o.jsx)(n.code,{children:"collectionMappings"})," registry;"]}),"\n",(0,o.jsx)(n.li,{children:"both collections and topics must be defined in the MongoDB cluster and the Kafka Cluster, respectively."}),"\n"]}),"\n",(0,o.jsx)(n.h2,{id:"messages-spec",children:"Messages Spec"}),"\n",(0,o.jsx)(n.p,{children:"Output Kafka messages key is compliant with the following schema:"}),"\n",(0,o.jsx)(n.pre,{children:(0,o.jsx)(n.code,{className:"language-json",children:'{\n "type": "object",\n "properties": {\n "$oid": {\
1n "type": "string"\n }\n }\n}\n'})}),"\n",(0,o.jsx)(n.p,{children:"Output Kafka messages payload is compliant the following schema:"}),"\n",(0,o.jsx)(n.pre,{children:(0,o.jsx)(n.code,{className:"language-json",children:'{\n "type": "object",\n "oneOf": [\n {\n "properties": {\n "op": {\n "const": "c",\n "description": "insert"\n },\n "before": { "type": "null" },\n "after": { "type": "object" }\n }\n },\n {\n "properties": {\n "op": {\n "const": "r",\n "description": "snapshot"\n },\n "before": { "type": "null" },\n "after": { "type": "object" }\n }\n },\n {\n "properties": {\n "op": {\n "const": "u",\n "description": "update"\n },\n "before": { "type": "object" },\n "after": { "type": "object" }\n }\n },\n {\n "properties": {\n "op": {\n "const": "d",\n "description": "delete"\n },\n "before": { "type": "object" },\n "after": { "type": "null" }\n }\n }\n ]\n}\n'})}),"\n",(0,o.jsx)(n.admonition,{type:"note",children:(0,o.jsxs)(n.p,{children:["Generated messages are compliant with ",(0,o.jsx)(n.a,{href:"/docs/products/fast_data_v2/concepts#fast-data-message-format",children:"Fast Data message format"})]})})]})}function p(e={}){let{wrapper:n}={...(0,a.R)(),...e.components};return n?(0,o.jsx)(n,{...e,children:(0,o.jsx)(l,{...e})}):l(e)}},201104(e,n,t){t.d(n,{A:()=>s});let s=t.p+"assets/images/mongezium_init_flow-dde38764e2c9801edbcbb18215fc6b2c.svg"},28453(e,n,t){t.d(n,{R:()=>i,x:()=>r});var s=t(296540);let o={},a=s.createContext(o);function i(e){let n=s.useContext(a);return s.useMemo(function(){return"function"==typeof e?e(n):{...n,...e}},[n,e])}function r(e){let n;return n=e.disableParentContext?"function"==typeof e.components?e.components(o):e.components||o:i(e.components),s.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.