PageSourceSearch

https://questdb.com/docs/assets/js/a4b7bc42.f7b19f68.js

js questdb.com collected 2026-10-02 05:06:13 UTC 69,967 bytes, 1 lines download raw bytes

1"use strict";(self.webpackChunkdocumentation=self.webpackChunkdocumentation||[]).push([["1754"],{2943:function(e,n,s){s.r(n),s.d(n,{frontMatter:()=>d,toc:()=>o,default:()=>h,metadata:()=>r,assets:()=>a,contentTitle:()=>c});var r=JSON.parse('{"id":"connect/message-brokers/kafka","title":"Ingest data from Kafka","description":"Stream data from Kafka to QuestDB. Set up the connector, map fields and timestamps, handle failed records, and prevent duplicates.","source":"@site/documentation/connect/message-brokers/kafka.md","sourceDirName":"connect/message-brokers","slug":"/connect/message-brokers/kafka","permalink":"/docs/connect/message-brokers/kafka","draft":false,"unlisted":false,"editUrl":"https://github.com/questdb/documentation/edit/main/documentation/connect/message-brokers/kafka.md","tags":[],"version":"current","frontMatter":{"slug":"/connect/message-brokers/kafka","title":"Ingest data from Kafka","sidebar_label":"Kafka","description":"Stream data from Kafka to QuestDB. Set up the connector, map fields and timestamps, handle failed records, and prevent duplicates."},"sidebar":"docs","previous":{"title":"Agents","permalink":"/docs/connect/agents"},"next":{"title":"Telegraf","permalink":"/docs/connect/message-brokers/telegraf"}}'),t=s(85893),i=s(50065);let d={slug:"/connect/message-brokers/kafka",title:"Ingest data from Kafka",sidebar_label:"Kafka",description:"Stream data from Kafka to QuestDB. Set up the connector, map fields and timestamps, handle failed records, and prevent duplicates."},c=void 0,a={},o=[{value:"Choosing an integration strategy",id:"choosing-an-integration-strategy",level:2},{value:"QuestDB Kafka connector",id:"questdb-kafka-connect-connector",level:2},{value:"Quick start",id:"quick-start",level:3},{value:"Prerequisites",id:"prerequisites",level:4},{value:"Step 1: Install the connector",id:"step-1-install-the-connector",level:4},{value:"Step 2: Configure the connector",id:"step-2-configure-the-connector",level:4},{value:"Step 3: Start the connector",id:"step-3-start-the-connector",level:4},{value:"Step 4: Test the pipeline",id:"step-4-test-the-pipeline",level:4},{value:"How data is mapped",id:"how-data-is-mapped",level:3},{value:"Designated timestamps",id:"designated-timestamps",level:3},{value:"Using a message field",id:"using-a-message-field",level:4},{value:"Using Kafka timestamps",id:"using-kafka-timestamps",level:4},{value:"Parsing string timestamps",id:"parsing-string-timestamps",level:4},{value:"Composed timestamps",id:"composed-timestamps",level:4},{value:"Type handling",id:"type-handling",level:3},{value:"Symbol columns",id:"symbol-columns",level:4},{value:"Numeric type inference",id:"numeric-type-inference",level:4},{value:"Target table options",id:"target-table-options",level:3},{value:"Table naming",id:"table-naming",level:4},{value:"Schema management",id:"schema-management",level:4},{value:"Delivery guarantees",id:"fault-tolerance",level:3},{value:"Exactly-once delivery",id:"exactly-once-delivery",level:4},{value:"Outages and reconnects",id:"outages-and-reconnects",level:4},{value:"Failover between QuestDB nodes",id:"failover-between-questdb-nodes",level:4},{value:"Dead letter queue",id:"dead-letter-queue",level:4},{value:"Shutdown and rebalances",id:"shutdown-and-rebalances",level:4},{value:"Performance tuning",id:"performance-tuning",level:3},{value:"Batch size and latency",id:"batch-size-and-latency",level:4},{value:"Backpressure",id:"backpressure",level:4},{value:"Raw JSON fast path",id:"raw-json-fast-path",level:4},{value:"Transformations",id:"transformations",level:3},{value:"OrderBookToArray",id:"orderbooktoarray",level:4},{value:"StructArrayExplode",id:"structarrayexplode",level:4},{value:"Legacy ILP transports",id:"legacy-ilp-transports",level:3},{value:"Configuration reference",id:"configuration-reference",level:3},{value:"Connector options",id:"connector-options",level:4},{value:"QWP delivery options",id:"qwp-delivery-options",level:4},{value:"Client configuration string",id:"client-configuration-string",level:4},{value:"Batching and buffer options",id:"batching-and-buffer-options",level:5}
1,{value:"Environment variable expansion",id:"environment-variable-expansion",level:5},{value:"Sample projects",id:"sample-projects",level:3},{value:"Stream processing",id:"stream-processing",level:2},{value:"Custom program",id:"custom-program",level:2},{value:"FAQ",id:"faq",level:2},{value:"See also",id:"see-also",level:2}];function l(e){let n={a:"a",admonition:"admonition",blockquote:"blockquote",code:"code",h2:"h2",h3:"h3",h4:"h4",h5:"h5",li:"li",ol:"ol",p:"p",pre:"pre",strong:"strong",table:"table",tbody:"tbody",td:"td",th:"th",thead:"thead",tr:"tr",ul:"ul",...(0,i.a)(),...e.components},{Details:s}=n;return s||function(e,n){throw Error("Expected "+(n?"component":"object")+" `"+e+"` to be defined: you likely forgot to import, pass, or provide it.")}("Details",!0),(0,t.jsxs)(t.Fragment,{children:[(0,t.jsxs)(n.p,{children:["Use the QuestDB Kafka connector to stream data from Apache Kafka into QuestDB\ntables. It handles data conversion, batching, and reconnects automatically.\nFollow the ",(0,t.jsx)(n.a,{href:"#quick-start",children:"quick start"})," to send your first message. If you have\nan existing pipeline that uses the InfluxDB Line Protocol (ILP) over HTTP, see\n",(0,t.jsx)(n.a,{href:"#legacy-ilp-transports",children:"migrating from ILP"}),"."]}),"\n",(0,t.jsx)(n.h2,{id:"choosing-an-integration-strategy",children:"Choosing an integration strategy"}),"\n",(0,t.jsx)(n.p,{children:"There are three ways to get data from Kafka into QuestDB:"}),"\n",(0,t.jsxs)(n.table,{children:[(0,t.jsx)(n.thead,{children:(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.th,{children:"Strategy"}),(0,t.jsx)(n.th,{children:"Recommended for"}),(0,t.jsx)(n.th,{children:"Complexity"})]})}),(0,t.jsxs)(n.tbody,{children:[(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:(0,t.jsx)(n.a,{href:"#questdb-kafka-connect-connector",children:"QuestDB Kafka connector"})}),(0,t.jsx)(n.td,{children:"Most users"}),(0,t.jsx)(n.td,{children:"Low"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:(0,t.jsx)(n.a,{href:"#stream-processing",children:"Stream processing (Flink)"})}),(0,t.jsx)(n.td,{children:"Complex transformations"}),(0,t.jsx)(n.td,{children:"Medium"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:(0,t.jsx)(n.a,{href:"#custom-program",children:"Custom program"})}),(0,t.jsx)(n.td,{children:"Special requirements"}),(0,t.jsx)(n.td,{children:"High"})]})]})]}),"\n",(0,t.jsx)(n.h2,{id:"questdb-kafka-connect-connector",children:"QuestDB Kafka connector"}),"\n",(0,t.jsxs)(n.p,{children:["The ",(0,t.jsx)(n.a,{href:"https://github.com/questdb/kafka-questdb-connector",children:"QuestDB Kafka connector"}),"\nis built on the ",(0,t.jsx)(n.a,{href:"https://docs.confluent.io/platform/current/connect/index.html",children:"Kafka Connect framework"}),".\nIt also works with Kafka-compatible systems such as\n",(0,t.jsx)(n.a,{href:"/docs/connect/message-brokers/redpanda/",children:"Redpanda"}),"."]}),"\n",(0,t.jsxs)(n.p,{children:["For new pipelines, use the\n",(0,t.jsx)(n.a,{href:"/docs/connect/wire-protocols/overview/",children:"QuestDB Wire Protocol (QWP)"})," with\n",(0,t.jsx)(n.code,{children:"ws::"})," or ",(0,t.jsx)(n.code,{children:"wss::"})," (TLS). The connector records progress in Kafka only after\nQuestDB confirms delivery. Enable ",(0,t.jsx)(n.a,{href:"#exactly-once-delivery",children:"deduplication"})," if\nyour table must not contain duplicate events."]}),"\n",(0,t.jsx)(n.h3,{id:"quick-start",children:"Quick start"}),"\n",(0,t.jsx)(n.p,{children:"This guide walks through setting up the connector to read JSON data from Kafka\nand write it to QuestDB."}),"\n",(0,t.jsx)(n.h4,{id:"prerequisites",children:"Prerequisites"}),"\n",(0,t.jsxs)(n.ul,{children:["\n",(0,t.jsx)(n.li,{children:"A running Apache Kafka 3.6 or newer broker (or a compatible system)"}),"\n",(0,t.jsx)(n.li,{children:"A running QuestDB 10.0 or newer instance, reachable on port 9000"}),"\n",(0,t.jsx)(n.li,{children:"QuestDB Kafka connector 0.24 or newer"}),"\n",(0,t.jsx)(n.li,{children:"Java 17+ (JDK)"}),"\n"]}),"\n",(0,t.jsxs)(n.p,{children:["The examples use Kafka at ",(0,t.jsx)(n.code,{children:"localhost:9092"})," and QuestDB at ",(0,t.jsx)(n.code,{children:"localhost:9000"}),".\nFor Kafka setup, follow the ",(0,t.jsx)(n.a,{href:"https://kafka.apache.org/quickstart/",children:"Apache Kafka quick start"}),"."]}),"\n",(0,t.jsx)(n.h4,{id:"step-1-install-the-connector",children:"Step 1: Install the connector"}),"\n",(0,t.jsxs)(n.p,{children:["Download the ",(0,t.jsx)(n.code,{children:"kafka-questdb-connector-<version>-bin.zip"})," archive from the\n",(0,t.jsx)(n.a,{href:"https://github.com/questdb/kafka-questdb-connector/releases",children:"connector releases"}),"."]}),"\n",(0,t.jsx)(n.p,{children:"Extract and copy to your Kafka installation:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-shell",children:"unzip kafka-questdb-connector-*-bin.zip\ncd kafka-questdb-connector\ncp ./*.jar /path/to/kafka_*.*-*.*.*/libs\n"})}),"\n",(0,t.jsx)(n.admonition,{type:"info",children:(0,t.jsxs)(n.p,{children:["The connector is also available from\n",(0,t.jsx)(n.a,{href:"https://www.confluent.io/hub/questdb/kafka-questdb-connector",children:"Confluent Hub"}),".\nFor Confluent platform users, see the\n",(0,t.jsx)(n.a,{href:"https://github.com/questdb/kafka-questdb-connector/tree/main/kafka-questdb-connector-samples/confluent-docker-images",children:"Confluent Docker images sample"}),"."]})}),"\n",(0,t.jsx)(n.h4,{id:"step-2-configure-the-connector",children:"Step 2: Configure the connector"}),"\n",(0,t.jsxs)(n.p,{children:["Create a configuration file at ",(0,t.jsx)(n.code,{children:"/path/to/kafka/config/questdb-connector.properties"}),":"]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",metastring:'title="questdb-connector.properties"',children:"name=questdb-sink\nconnector.class=io.questdb.kafka.QuestDBSinkConnector\n\n# QuestDB connection (QWP over WebSocket)\nclient.conf.string=ws::addr=localhost:9000;\n\n# Kafka source\ntopics=example-topic\n\n# Target table (optional - defaults to topic name)\ntable=example_table\n\n# Message format\nkey.converter=org.apache.kafka.connect.storage.StringConverter\nvalue.converter=org.apache.kafka.connect.json.JsonConverter\nvalue.converter.sc
1hemas.enable=false\ninclude.key=false\n"})}),"\n",(0,t.jsx)(n.h4,{id:"step-3-start-the-connector",children:"Step 3: Start the connector"}),"\n",(0,t.jsxs)(n.p,{children:["In Kafka's ",(0,t.jsx)(n.code,{children:"config/connect-standalone.properties"}),", set ",(0,t.jsx)(n.code,{children:"bootstrap.servers"})," to\nyour broker address (",(0,t.jsx)(n.code,{children:"localhost:9092"})," for this example). From your Kafka\ninstallation directory, create the topic and start the connector:"]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-shell",children:"bin/kafka-topics.sh --create --if-not-exists --topic example-topic --bootstrap-server localhost:9092\nbin/connect-standalone.sh config/connect-standalone.properties config/questdb-connector.properties\n"})}),"\n",(0,t.jsx)(n.h4,{id:"step-4-test-the-pipeline",children:"Step 4: Test the pipeline"}),"\n",(0,t.jsx)(n.p,{children:"From another terminal in your Kafka installation directory, publish a test message:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-shell",children:"bin/kafka-console-producer.sh --topic example-topic --bootstrap-server localhost:9092\n"})}),"\n",(0,t.jsx)(n.p,{children:"Enter this JSON (as a single line):"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-json",children:'{"symbol": "AAPL", "price": 192.34, "volume": 1200}\n'})}),"\n",(0,t.jsx)(n.p,{children:"Verify the data in QuestDB:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-shell",children:"curl -G --data-urlencode \"query=select * from 'example_table'\" http://localhost:9000/exp\n"})}),"\n",(0,t.jsx)(n.p,{children:"Expected output:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-csv",children:'"symbol","price","volume","timestamp"\n"AAPL",192.34,1200,"2026-02-03T15:10:00.000000Z"\n'})}),"\n",(0,t.jsx)(n.p,{children:"QuestDB assigns the timestamp when the message arrives, so your value will differ."}),"\n",(0,t.jsxs)(n.p,{children:["Next, choose a ",(0,t.jsx)(n.a,{href:"#designated-timestamps",children:"timestamp source"})," and review\n",(0,t.jsx)(n.a,{href:"#fault-tolerance",children:"delivery guarantees"})," before using the connector in production."]}),"\n",(0,t.jsx)(n.h3,{id:"how-data-is-mapped",children:"How data is mapped"}),"\n",(0,t.jsx)(n.p,{children:"The connector converts each Kafka message field to a QuestDB column. Nested\nstructures and maps are flattened with underscores."}),"\n",(0,t.jsx)(n.p,{children:(0,t.jsx)(n.strong,{children:"Example input:"})}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-json",children:'{\n  "firstname": "John",\n  "lastname": "Doe",\n  "age": 30,\n  "address": {\n    "street": "Main Street",\n    "city": "New York"\n  }\n}\n'})}),"\n",(0,t.jsx)(n.p,{children:(0,t.jsx)(n.strong,{children:"Resulting table:"})}),"\n",(0,t.jsxs)(n.table,{children:[(0,t.jsx)(n.thead,{children:(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.th,{children:"firstname"}),(0,t.jsx)(n.th,{children:"lastname"}),(0,t.jsx)(n.th,{children:"age"}),(0,t.jsx)(n.th,{children:"address_street"}),(0,t.jsx)(n.th,{children:"address_city"})]})}),(0,t.jsx)(n.tbody,{children:(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"John"}),(0,t.jsx)(n.td,{children:"Doe"}),(0,t.jsx)(n.td,{children:"30"}),(0,t.jsx)(n.td,{children:"Main Street"}),(0,t.jsx)(n.td,{children:"New York"})]})})]}),"\n",(0,t.jsx)(n.h3,{id:"designated-timestamps",children:"Designated timestamps"}),"\n",(0,t.jsxs)(n.p,{children:["The connector supports four strategies for\n",(0,t.jsx)(n.a,{href:"/docs/concepts/designated-timestamp/",children:"designated timestamps"}),":"]}),"\n",(0,t.jsxs)(n.table,{children:[(0,t.jsx)(n.thead,{children:(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.th,{children:"Strategy"}),(0,t.jsx)(n.th,{children:"Configuration"}),(0,t.jsx)(n.th,{children:"Use case"})]})}),(0,t.jsxs)(n.tbody,{children:[(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"Server-assigned"}),(0,t.jsx)(n.td,{children:"(default)"}),(0,t.jsx)(n.td,{children:"QuestDB assigns timestamp on receipt"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"Message payload"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"timestamp.field.name=fieldname"})}),(0,t.jsx)(n.td,{children:"Use a field from the message"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"Kafka metadata"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"timestamp.kafk
1a.native=true"})}),(0,t.jsx)(n.td,{children:"Use Kafka's message timestamp"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"Composed"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"timestamp.field.name=date,time"})}),(0,t.jsx)(n.td,{children:"Combine multiple fields"})]})]})]}),"\n",(0,t.jsx)(n.p,{children:"These strategies are mutually exclusive."}),"\n",(0,t.jsx)(n.h4,{id:"using-a-message-field",children:"Using a message field"}),"\n",(0,t.jsx)(n.p,{children:"If your message contains a timestamp field:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"timestamp.field.name=event_time\ntimestamp.units=millis\n"})}),"\n",(0,t.jsxs)(n.p,{children:["Supported units are ",(0,t.jsx)(n.code,{children:"nanos"}),", ",(0,t.jsx)(n.code,{children:"micros"}),", ",(0,t.jsx)(n.code,{children:"millis"}),", ",(0,t.jsx)(n.code,{children:"seconds"}),", and ",(0,t.jsx)(n.code,{children:"auto"})," (the\ndefault). Auto-detection supports timestamps after April 26, 1970."]}),"\n",(0,t.jsx)(n.h4,{id:"using-kafka-timestamps",children:"Using Kafka timestamps"}),"\n",(0,t.jsx)(n.p,{children:"To use Kafka's built-in message timestamp:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"timestamp.kafka.native=true\n"})}),"\n",(0,t.jsx)(n.h4,{id:"parsing-string-timestamps",children:"Parsing string timestamps"}),"\n",(0,t.jsx)(n.p,{children:"For timestamps stored as strings:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"timestamp.field.name=created_at\ntimestamp.string.fields=updated_at,deleted_at\ntimestamp.string.format=yyyy-MM-dd HH:mm:ss.SSSUUU z\n"})}),"\n",(0,t.jsxs)(n.p,{children:["The ",(0,t.jsx)(n.code,{children:"timestamp.field.name"})," field becomes the designated timestamp. Fields in\n",(0,t.jsx)(n.code,{children:"timestamp.string.fields"})," are parsed as regular timestamp columns."]}),"\n",(0,t.jsxs)(n.p,{children:["See ",(0,t.jsx)(n.a,{href:"/docs/query/functions/date-time/#timestamp-format",children:"QuestDB timestamp format"}),"\nfor format patterns."]}),"\n",(0,t.jsx)(n.h4,{id:"composed-timestamps",children:"Composed timestamps"}),"\n",(0,t.jsx)(n.p,{children:"Some data sources split timestamps across multiple fields (common with KDB-style data):"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-json",children:'{\n  "symbol": "BTC-USD",\n  "date": "20260202",\n  "time": "135010207"\n}\n'})}),"\n",(0,t.jsx)(n.p,{children:"Configure the connector to concatenate and parse them:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"timestamp.field.name=date,time\ntimestamp.string.format=yyyyMMddHHmmssSSS\n"})}),"\n",(0,t.jsxs)(n.p,{children:["The fields ",(0,t.jsx)(n.code,{children:"date"})," and ",(0,t.jsx)(n.code,{children:"time"})," are concatenated into ",(0,t.jsx)(n.code,{children:"20260202135010207"}),", parsed\nto produce ",(0,t.jsx)(n.code,{children:"2026-02-02T13:50:10.207000Z"}),". The source fields are consumed and do\nnot appear as columns in the output."]}),"\n",(0,t.jsx)(n.p,{children:"All listed fields must be present in each message."}),"\n",(0,t.jsx)(n.h3,{id:"type-handling",children:"Type handling"}),"\n",(0,t.jsx)(n.h4,{id:"symbol-columns",children:"Symbol columns"}),"\n",(0,t.jsxs)(n.p,{children:["Use the ",(0,t.jsx)(n.code,{children:"symbols"})," option to create columns as\n",(0,t.jsx)(n.a,{href:"/docs/concepts/symbol/",children:"symbol"})," type for better performance on\nrepeated string values:"]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"symbols=instrument,exchange,currency\n"})}),"\n",(0,t.jsx)(n.h4,{id:"numeric-type-inference",children:"Numeric type inference"}),"\n",(0,t.jsx)(n.p,{children:"Without a schema, the connector infers types from values. This can cause issues\nwhen a field is sometimes an integer and sometimes a float:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-json",children:'{"volume": 42}      // Inferred as long\n{"volume": 42.5}    // Error: column is long, value is double\n'})}),"\n",(0,t.jsx)(n.p,{children:"Solutions:"}),"\n",(0,t.jsxs)(n.ol,{children:["\n",(0,t.jsxs)(n.li,{children:["Use the ",(0,t.jsx)(n.code,{children:"doubles"})," option to force double type:","\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"doubles=volume,price\n"})}),"\n"]}),"\n",(0,t.jsxs)(n.li,{children:["Pre-create the table with explicit column types using\n",(0,t.jsx)(n.a,{href:"/docs/query/sql/create-table/",children:"CREATE TABLE"})]}),"\n"]}),"\n",(0,t.jsx)(n.h3,{id:"target-table-options",children:"Target table options"}),"\n",(0,t.jsx)(n.h4,{id:"table-naming",children:"Table naming"}),"\n",(0,t.jsx)(n.p,{children:"By default, the table name matches the Kafka topic name. Override with:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"table=my_custom_table\n"})}),"\n",(0,t.jsxs)(n.p,{children:["The ",(0,t.jsx)(n.code,{children:"table"})," option supports templating:"]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"table=kafka_${topic}_${partition}\n"})}),"\n",(0,t.jsxs)(n.p,{children:["Available variables: ",(0,t.jsx)(n.code,{children:"${topic}"}),", ",(0,t.jsx)(n.code,{children:"${key}"}),", ",(0,t.jsx)(n.code,{children:"${partition}"})]}),"\n",(0,t.jsxs)(n.p,{children:["If ",(0,t.jsx)(n.code,{children:"${key}"})," is used and the message has no key, it resolves to ",(0,t.jsx)(n.code,{children:"null"}),"."]}),"\n",(0,t.jsx)(n.h4,{id:"schema-management",children:"Schema management"}),"\n",(0,t.jsxs)(n.p,{children:["Tables are created automatically when they don't exist. This is convenient for\ndevelopment but in production, pre-create tables using\n",(0,t.jsx)(n.a,{href:"/docs/query/sql/create-table/",children:"CREATE TABLE"})," for control over partitioning,\nindexes, and column types."]}),"\n",(0,t.jsx)(n.h3,{id:"fault-tolerance",children:"Delivery guarantees"}),"\n",(0,t.jsxs)(n.p,{children:["With QWP, delivery is ",(0,t.jsx)(n.strong,{children:"at least once"}),": the connector commits Kafka offsets\nonly after QuestDB acknowledges the corresponding rows. Records without\nconfirmation remain uncommitted and can be retried. You can route invalid\nrecords to a ",(0,t.jsx)(n.a,{href:"#dead-letter-queue",children:"dead letter queue"}),"."]}),"\n",(0,t.jsxs)(n.p,{children:["Retries can produce duplicates. For example, QuestDB may save a row just\nbefore the connection drops, leaving the connector unsure whether it arrived.\nUse deduplication if each event must appear only once, and keep source records\nin Kafka long enough for ",(0,t.jsx)(n.a,{href:"#outages-and-reconnects",children:"outage recovery"}),"."]}),"\n",(0,t.jsx)(n.h4,{id:"exactly-once-delivery",children:"Exactly-once delivery"}),"\n",(0,t.jsxs)(n.p,{children:["For exactly-once results, enable\n",(0,t.jsx)(n.a,{href:"/docs/concepts/deduplication/",children:"deduplication"})," on the target table with keys\nthat identify a unique event:"]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-questdb-sql",children:"CREATE TABLE trades (\n    timestamp TIMESTAMP,\n    trade_id LONG,\n    symbol SYMBOL,\n    price DOUBLE,\n    volume LONG\n) TIMESTAMP(timestamp) PARTITION BY DAY\nDEDUP UPSERT KEYS(timestamp, trade_id);\n"})}),"\n",(0,t.jsxs)(n.p,{children:["Here, ",(0,t.jsx)(n.code,{children:"trade_id"})," is an event identifier supplied by your producer. Choose keys\nthat distinguish separate events, even when they share a timestamp."]}),"\n",(0,t.jsxs)(n.p,{children:["Use a timestamp from the ",(0,t.jsx)(n.a,{href:"#using-a-message-field",children:"message payload"})," or\n",(0,t.jsx)(n.a,{href:"#using-kafka-timestamps",children:"Kafka metadata"})," so it stays the same on retry.\nThe default server-assigned timestamp changes on retry and cannot deduplicate\nthe event. See\n",(0,t.jsx)(n.a,{href:"/docs/concepts/delivery-semantics/",children:"Delivery semantics"})," for the full model."]}),"\n",(0,t.jsx)(n.h4,{id:"outages-and-reconnects",children:"Outages and reconnects"}),"\n",(0,t.jsxs)(n.p,{children:["The connector reconnects and retries automatically after a connection drops.\nIf QuestDB is unreachable when a task starts, it retries every\n",(0,t.jsx)(n.code,{children:"retry.backoff.ms"})," (default 3 seconds). Authentication and configuration errors\nfail the task immediately; fix the error before restarting it."]}),"\n",(0,t.jsxs)(n.p,{children:["If rows are pending and QuestDB confirms n
1o further delivery for\n",(0,t.jsx)(n.code,{children:"qwp.progress.timeout.ms"})," (default 5 minutes), the task fails. Restart it once\nQuestDB is available. Increase this timeout if you need to tolerate longer\noutages. A backlog that continues to drain resets the timer."]}),"\n",(0,t.jsxs)(n.p,{children:[(0,t.jsx)(n.strong,{children:"Set Kafka retention to cover the outage and catch-up time."})," A restarted\ntask resumes from its last committed offset only if those records still\nexist. Kafka can expire uncommitted records; expired records cannot be\nrecovered by the connector."]}),"\n",(0,t.jsx)(n.h4,{id:"failover-between-questdb-nodes",children:"Failover between QuestDB nodes"}),"\n",(0,t.jsxs)(n.p,{children:["With QuestDB Enterprise, list every node of the cluster in ",(0,t.jsx)(n.code,{children:"addr"})," and the\nconnector follows whichever node holds the primary role:"]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"client.conf.string=wss::addr=node-a:9000,node-b:9000;token=${QUESTDB_TOKEN};\n"})}),"\n",(0,t.jsxs)(n.p,{children:["Failover itself is a manual operation, see\n",(0,t.jsx)(n.a,{href:"/docs/high-availability/failover/",children:"Failover and role switch"}),". Once you promote\na replica, the connector switches to it on its own:"]}),"\n",(0,t.jsxs)(n.ol,{children:["\n",(0,t.jsx)(n.li,{children:"The demoted node closes the connection on the first write it refuses."}),"\n",(0,t.jsx)(n.li,{children:"The client reconnects and tries each address in turn. Replicas reject the\nconnection immediately, so it lands on the new primary without a backoff\ndelay."}),"\n",(0,t.jsx)(n.li,{children:"Rows that were not acknowledged before the switch are re-sent to the new\nprimary."}),"\n"]}),"\n",(0,t.jsxs)(n.p,{children:["Step 3 can produce duplicates, so enable\n",(0,t.jsx)(n.a,{href:"#exactly-once-delivery",children:"deduplication"})," on the affected tables."]}),"\n",(0,t.jsxs)(n.p,{children:["While no node accepts writes, the client keeps retrying in the background. The\ntask fails only if ",(0,t.jsx)(n.code,{children:"qwp.progress.timeout.ms"})," elapses without progress, and a\ntask that starts during the switch retries every ",(0,t.jsx)(n.code,{children:"retry.backoff.ms"}),"."]}),"\n",(0,t.jsx)(n.h4,{id:"dead-letter-queue",children:"Dead letter queue"}),"\n",(0,t.jsxs)(n.p,{children:["Configure a dead letter queue (DLQ) to set aside invalid records for inspection\nwhile valid records continue to QuestDB. Add these settings to your connector\nconfiguration (",(0,t.jsx)(n.code,{children:"questdb-connector.properties"}),", or the connector JSON in\ndistributed mode):"]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",metastring:'title="questdb-connector.properties"',children:"errors.tolerance=all\nerrors.deadletterqueue.topic.name=dlq-questdb\n# Use 1 for a single-broker development cluster\nerrors.deadletterqueue.topic.replication.factor=1\n"})}),"\n",(0,t.jsxs)(n.p,{children:["Both ",(0,t.jsx)(n.code,{children:"errors.tolerance=all"})," and a DLQ topic are required. Choose a replication\nfactor appropriate for your production Kafka cluster."]}),"\n",(0,t.jsxs)(n.p,{children:["By default, the connector sends records with conversion errors, oversized rows,\nor schema mismatches to the DLQ. For example, a string sent to a ",(0,t.jsx)(n.code,{children:"DOUBLE"}),"\ncolumn is a schema mismatch. Without a usable DLQ, these errors stop the task.\nAuthentication errors and other server failures still stop the task."]}),"\n",(0,t.jsxs)(n.p,{children:["The connector retries rejected batches to identify the invalid records. This\ncan slow ingestion. Set ",(0,t.jsx)(n.code,{children:"dlq.send.batch.on.error=true"})," only if you prefer to\nsend the entire rejected batch to the DLQ, including any valid records in it."]}),"\n",(0,t.jsxs)(n.p,{children:["A common cause of schema mismatches is a JSON field that switches between\ninteger and float. Pin such fields with the ",(0,t.jsx)(n.code,{children:"doubles"})," option or pre-create the table, see\n",(0,t.jsx)(n.a,{href:"#numeric-type-inference",children:"Numeric type inference"}),"."]}),"\n",(0,t.jsxs)(n.p,{children:["See the ",(0,t.jsx)(n.a,{href:"https://developer.confluent.io/courses/kafka-connect/error-handling-and-dead-letter-queues/",children:"Confluent DLQ documentation"}),"\nfor details."]}),"\n",(0,t.jsx)(n.h4,{id:"shutdown-and-rebalances",children:"Shutdown and rebalances"}),"\n",(0,t.jsxs)(n.p,{children:["During a normal shutdown or rebalance, the connector sends pending rows and\nwaits up to ",(0,t.jsx)(n.code,{children:"qwp.commit.ack.timeout.ms"})," (default 500 ms) for confirmation before\nKafka Connect commits offsets. Unconfirmed records remain uncommitted and may\nbe delivered again by the next task. Deduplication prevents these retries\nfrom creating duplicate rows."]}),"\n",(0,t.jsx)(n.h3,{id:"performance-tuning",children:"Performance tuning"}),"\n",(0,t.jsx)(n.p,{children:"Start with the defaults. Adjust batching if messages take too long to appear,\nor buffer limits if network latency keeps the connector waiting for delivery\nconfirmations."}),"\n",(0,t.jsx)(n.h4,{id:"batch-size-and-latency",children:"Batch size and latency"}),"\n",(0,t.jsxs)(n.table,{children:[(0,t.jsx)(n.thead,{children:(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.th,{children:"Setting"}),(0,t.jsx)(n.th,{children:"Default"}),(0,t.jsx)(n.th,{children:"Use it to"})]})}),(0,t.jsxs)(n.tbody,{children:[(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"auto_flush_rows"})}),(0,t.jsx)(n.td,{children:"75000 rows"}
1),(0,t.jsx)(n.td,{children:"Send a batch when it reaches this size"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"auto_flush_interval"})}),(0,t.jsx)(n.td,{children:"1000 ms"}),(0,t.jsx)(n.td,{children:"Send pending rows periodically, even while new records keep arriving"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"allowed.lag"})}),(0,t.jsx)(n.td,{children:"1000 ms"}),(0,t.jsx)(n.td,{children:"Limit how long the connector waits for more records before sending a partial batch when idle"})]})]})]}),"\n",(0,t.jsx)(n.p,{children:"For smaller batches and more frequent sends:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"client.conf.string=ws::addr=localhost:9000;auto_flush_rows=1000;auto_flush_interval=250;\nallowed.lag=250\n"})}),"\n",(0,t.jsx)(n.p,{children:"Smaller batches increase request overhead. The connector also sends pending\nrows when Kafka Connect commits offsets."}),"\n",(0,t.jsx)(n.h4,{id:"backpressure",children:"Backpressure"}),"\n",(0,t.jsx)(n.p,{children:"The connector automatically pauses consumption when too many rows are waiting\nfor QuestDB to confirm delivery. It resumes when QuestDB catches up."}),"\n",(0,t.jsxs)(n.ul,{children:["\n",(0,t.jsxs)(n.li,{children:[(0,t.jsx)(n.code,{children:"qwp.max.inflight.rows"})," (default 150,000) limits buffered and sent rows\nawaiting confirmation. This is a soft limit: a Kafka poll batch can exceed it."]}),"\n",(0,t.jsxs)(n.li,{children:[(0,t.jsx)(n.code,{children:"sf_max_total_bytes"})," (default 128 MiB) limits memory used to buffer encoded\nrows awaiting confirmation."]}),"\n"]}),"\n",(0,t.jsx)(n.p,{children:"On high-latency connections, larger limits allow more data to be sent while\nwaiting for confirmations, at the cost of more memory. For example:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"qwp.max.inflight.rows=500000\nclient.conf.string=ws::addr=questdb.example.com:9000;sf_max_total_bytes=512m;\n"})}),"\n",(0,t.jsxs)(n.p,{children:["If the byte buffer stays full for ",(0,t.jsx)(n.code,{children:"sf_append_deadline_millis"})," (default 30\nseconds), the connector reconnects and retries unconfirmed records from Kafka."]}),"\n",(0,t.jsx)(n.h4,{id:"raw-json-fast-path",children:"Raw JSON fast path"}),"\n",(0,t.jsx)(n.admonition,{title:"Experimental",type:"caution",children:(0,t.jsx)(n.p,{children:"Test this mode with representative messages before using it in production.\nKeep the default converter-based mode if you need value transformations or\nschema-defined column types."})}),"\n",(0,t.jsx)(n.p,{children:"For JSON object messages, you can let the connector parse the values directly\nto reduce conversion work:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"value.converter=org.apache.kafka.connect.converters.ByteArrayConverter\nvalue.format=json\n"})}),"\n",(0,t.jsxs)(n.p,{children:["For messages wrapped as ",(0,t.jsx)(n.code,{children:'{"schema": {...}, "payload": {...}}'}),", set\n",(0,t.jsx)(n.code,{children:"value.format=json_envelope"}),". Only ",(0,t.jsx)(n.code,{children:"payload"})," becomes the row; the schema is\nignored. Choose the mode explicitly: ",(0,t.jsx)(n.code,{children:"json"})," would turn the envelope into\n",(0,t.jsx)(n.code,{children:"schema_*"})," and ",(0,t.jsx)(n.code,{children:"payload_*"})," columns."]}),"\n",(0,t.jsx)(n.p,{children:"Before switching, check that:"}),"\n",(0,t.jsxs)(n.ul,{children:["\n",(0,t.jsx)(n.li,{children:"Your messages contain JSON objects. Top-level strings, numbers, and arrays\nare not supported."}),"\n",(0,t.jsxs)(n.li,{children:["You do not use transformations that read or modify message values, such as\nthe ",(0,t.jsx)(n.a,{href:"#transformations",children:"array transforms"}),". Topic routing transforms\nsuch as ",(0,t.jsx)(n.code,{children:"RegexRouter"})," still work."]}),"\n",(0,t.jsxs)(n.li,{children:["You do not use ",(0,t.jsx)(n.a,{href:"#composed-timestamps",children:"composed timestamps"}),"."]}),"\n",(0,t.jsxs)(n.li,{children:["Your table accepts types inferred from JSON values. Schema declarations\nsuch as ",(0,t.jsx)(n.code,{children:"INT8"})," or ",(0,t.jsx)(n.code,{children:"FLOAT32"})," are ignored; numbers become ",(0,t.jsx)(n.code,{children:"LONG"})," or ",(0,t.jsx)(n.code,{children:"DOUBLE"}),".\nUse ",(0,t.jsx)(n.code,{children:"doubles"})," for fields that must always be sent as doubles."]}),"\n"]}),"\n",(0,t.jsxs)(n.p,{children:["The key still uses ",(0,t.jsx)(n.code,{children:"key.converter"}),". Field mapping, nested-object flattening,\nand numeric arrays remain available."]}),"\n",(0,t.jsxs)(s,{children:[(0,t.jsx)("summary",{children:"Additional JSON compatibility details"}),(0,t.jsxs)(n.table,{children:[(0,t.jsx)(n.thead,{children:(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.th,{children:"Input or setting"}),(0,t.jsx)(n.th,{children:"Behavior in raw JSON mode"})]})}),(0,t.jsxs)(n.tbody,{children:[(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"Duplicate field names"}),(0,t.jsx)(n.td,{children:"The first value is kept; the standard converter keeps the last"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"Integers outside the signed 64-bit range"}),(0,t.jsx)(n.td,{children:"Written as doubles, which can lose precision"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"Empty field names"}),(0,t.jsxs)(n.td,{children:["Sent to the DLQ, or fail the task without one; the standard converter uses a column named ",(0,t.jsx)(n.code,{children:"value"})]})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsxs)(n.td,{children:["Objects or arrays listed in ",(0,t.jsx)(n.code,{children:"symbols"})]}),(0,t.jsx)(n.td,{children:"Remain flattened objects or arrays, rather than becoming symbol columns"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"Objects inside arrays"}),(0,t.jsxs)(n.td,{children:["Fail, or are skipped with ",(0,t.jsx)(n.code,{children:"skip.unsupported.types=true"})]})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"Nesting deeper than 64 levels"}
1),(0,t.jsx)(n.td,{children:"Rejected as invalid data"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"Auto-created tables"}),(0,t.jsx)(n.td,{children:"Column order follows the JSON document and may differ from converter-based ingestion"})]})]})]}),(0,t.jsxs)(n.p,{children:["Use ",(0,t.jsx)(n.code,{children:"ws"})," or ",(0,t.jsx)(n.code,{children:"http"})," (or their TLS variants). With the legacy TCP transport,\nmalformed JSON fails the task even if a DLQ is configured."]})]}),"\n",(0,t.jsx)(n.h3,{id:"transformations",children:"Transformations"}),"\n",(0,t.jsx)(n.h4,{id:"orderbooktoarray",children:"OrderBookToArray"}),"\n",(0,t.jsxs)(n.p,{children:["The connector includes an ",(0,t.jsx)(n.code,{children:"OrderBookToArray"})," Single Message Transform (SMT)\nfor converting arrays of structs into arrays of arrays. This is useful for\norder book data or tabular data stored as rows that needs to be pivoted into\ncolumnar form."]}),"\n",(0,t.jsxs)(n.p,{children:["For querying order book data stored as arrays, see\n",(0,t.jsx)(n.a,{href:"/docs/tutorials/order-book/",children:"Order book analytics using arrays"}),"."]}),"\n",(0,t.jsx)(n.p,{children:(0,t.jsx)(n.strong,{children:"Input:"})}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-json",children:'{\n  "symbol": "BTC-USD",\n  "buy_entries": [\n    { "price": 100.5, "size": 10.0 },\n    { "price": 99.8, "size": 25.0 }\n  ]\n}\n'})}),"\n",(0,t.jsx)(n.p,{children:(0,t.jsx)(n.strong,{children:"Output:"})}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-json",children:'{\n  "symbol": "BTC-USD",\n  "bids": [\n    [100.5, 99.8],\n    [10.0, 25.0]\n  ]\n}\n'})}),"\n",(0,t.jsx)(n.p,{children:(0,t.jsx)(n.strong,{children:"Configuration:"})}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"transforms=orderbook\ntransforms.orderbook.type=io.questdb.kafka.OrderBookToArray$Value\ntransforms.orderbook.mappings=buy_entries:bids:price,size;sell_entries:asks:price,size\n"})}),"\n",(0,t.jsxs)(n.p,{children:["The ",(0,t.jsx)(n.code,{children:"mappings"})," format is ",(0,t.jsx)(n.code,{children:"sourceField:targetField:field1,field2;..."})]}),"\n",(0,t.jsx)(n.p,{children:(0,t.jsx)(n.strong,{children:"Behavior:"})}),"\n",(0,t.jsxs)(n.ul,{children:["\n",(0,t.jsxs)(n.li,{children:["All extracted values are converted to ",(0,t.jsx)(n.code,{children:"double"})]}),"\n",(0,t.jsx)(n.li,{children:"Missing source fields are skipped (no error)"}),"\n",(0,t.jsx)(n.li,{children:"Empty source arrays are skipped"}),"\n",(0,t.jsx)(n.li,{children:"Null values inside structs cause an error"}),"\n",(0,t.jsx)(n.li,{children:"If the target field name already exists in the input, it is replaced"}),"\n",(0,t.jsx)(n.li,{children:"Works with both schema-based and schemaless messages"}),"\n"]}),"\n",(0,t.jsx)(n.admonition,{type:"note",children:(0,t.jsx)(n.p,{children:"QuestDB requires all inner arrays to have the same length. The OrderBookToArray\nSMT satisfies this naturally since each inner array comes from the same source\nentries."})}),"\n",(0,t.jsx)(n.h4,{id:"structarrayexplode",children:"StructArrayExplode"}),"\n",(0,t.jsxs)(n.p,{children:["The ",(0,t.jsx)(n.code,{children:"StructArrayExplode"})," SMT converts arrays of structs into ",(0,t.jsxs)(n.strong,{children:["separate 1D\n",(0,t.jsx)(n.code,{children:"double[]"})," columns"]}),", one per struct field. Unlike ",(0,t.jsx)(n.code,{children:"OrderBookToArray"}),' which\nproduces a single 2D array column, this transform "explodes" each struct field\ninto its own column.']}),"\n",(0,t.jsx)(n.p,{children:(0,t.jsx)(n.strong,{children:"Input:"})}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-json",children:'{\n  "symbol": "AAPL",\n  "vols": [\n    { "strike": 150.0, "ivol": 0.25 },\n    { "strike": 160.0, "ivol": 0.22 }\n  ]\n}\n'})}),"\n",(0,t.jsx)(n.p,{children:(0,t.jsx)(n.strong,{children:"Output:"})}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-json",children:'{\n  "symbol": "AAPL",\n  "strikes": [150.0, 160.0],\n  "ivols": [0.25, 0.22]\n}\n'})}),"\n",(0,t.jsx)(n.p,{children:(0,t.jsx)(n.strong,{children:"Configuration:"})}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"transforms=explode\ntransforms.explode.type=io.questdb.kafka.StructArrayExplode$Value\ntransforms.explode.mappings=vols:strikes,ivols:strike,ivol\n"})}),"\n",(0,t.jsxs)(n.p,{children:["The ",(0,t.jsx)(n.code,{children:"mappings"})," format is ",(0,t.jsx)(n.code,{children:"sourceField:targetCol1,targetCol2:structField1,structField2;..."})]}),"\n",(0,t.jsxs)(n.p,{children:["Target columns and struct fields are paired positionally: ",(0,t.jsx)(n.code,{children:"structField1"})," maps to\n",(0,t.jsx)(n.code,{children:"targetCol1"}),", ",(0,t.jsx)(n.code,{children:"structField2"})," maps to ",(0,t.jsx)(n.code,{children:"targetCol2"}),", and so on. The number of\ntarget columns must equal the number of struct fields."]}),"\n",(0,t.jsx)(n.p,{children:"Use semicolons to separate mappings from different source arrays:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"transforms.explode.mappings=bids:bid_prices,bid_amounts:price,amount;asks:ask_prices,ask_amounts:price,amount\n"})}),"\n",(0,t.jsx)(n.p,{children:(0,t.jsx)(n.strong,{children:"Behavior:"})}),"\n",(0,t.jsxs)(n.ul,{children:["\n",(0,t.jsxs)(n.li,{children:["All extracted values are converted to ",(0,t.jsx)(n.code,{children:"double"})]}),"\n",(0,t.jsx)(n.li,{children:"Missing source fields are skipped (no error)"}),"\n",(0,t.jsx)(n.li,{children:"Empty source arrays are skipped"}),"\n",(0,t.jsx)(n.li,{children:"Null values inside structs cause an error"}),"\n",(0,t.jsx)(n.li,{children:"If a target column name already exists in the input, it is replaced"}),"\n",(0,t.jsx)(n.li,{children:"Works with both schema-based and schemaless messages"}),"\n"]}),"\n",(0,t.jsx)(n.p,{children:(0,t.jsx)(n.strong,{children:"Comparison with OrderBookToArray:"})}),"\n",(0,t.jsxs)(n.table,{children:[(0,t.jsx)(n.thead,{children:(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.th,{}),(0,t.jsx)(n.th,{children:"OrderBookToArray"}),(0,t.jsx)(n.th,{children:"StructArrayExplode"})]})}),(0,t.jsxs)(n.tbody,{children:[(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"Output"}),(0,t.jsxs)(n.td,{children:["One 2D ",(0,t.jsx)(n.code,{children:"double[][]"})," column"]}),(0,t.jsxs)(n.td,{children:["Separate 1D ",(0,t.jsx)(n.code,{children:"double[]"})," columns"]})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"Mapping format"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"source:target:field1,field2"})}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"source:target1,target2:field1,field2"})})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"Use case"}),(0,t.jsx)(n.td,{children:"All fields in one array column"}),(0,t.jsx)(n.td,{children:"Each field as its own column"})]})]})]}),"\n",(0,t.jsx)(n.h3,{id:"legacy-ilp-transports",children:"Legacy ILP transports"}),"\n",(0,t.jsxs)(n.p,{children:["The ",(0,t.jsx)(n.code,{children:"http"})," and ",(0,t.jsx)(n.code,{children:"https"})," protocols send rows as\n",(0,t.jsx)(n.a,{href:"/docs/connect/compatibility/ilp/overview/",children:"InfluxDB Line Protocol"})," over HTTP.\nUse them with QuestDB versions before 10.0 or to keep an existing HTTP\npipeline. For new pipelines on QuestDB 10.0 or newer, use ",(0,t.jsx)(n.code,{children:"ws"})," or ",(0,t.jsx)(n.code,{children:"wss"}),"."]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"client.conf.string=http::addr=localhost:9000;retry_timeout=60000;\n"})}),"\n",(0,t.jsxs)(n.p,{children:["HTTP retries temporary errors for up to ",(0,t.jsx)(n.code,{children:"retry_timeout"})," milliseconds (default\n10,000), then fails the task. The ",(0,t.jsx)(n.code,{children:"qwp.*"})," options have no effect. You can still\nuse a DLQ for invalid records and deduplication to prevent duplicates on retry."]}),"\n",(0,t.jsx)(n.p,{children:"To migrate an HTTP pipeline to QWP:"}),"\n",(0,t.jsxs)(n.ol,{children:["\n",(0,t.jsx)(n.li,{children:"Upgrade to QuestDB 10.0 or newer and connector 0.24 or newer."}),"\n",(0,t.jsxs)(n.li,{children:["Change ",(0,t.jsx)(n.code,{children:"http::"})," to ",(0,t.jsx)(n.code,{children:"ws::"}),", or ",(0,t.jsx)(n.code,{children:"https::"})," to ",(0,t.jsx)(n.code,{children:"wss::"})," for TLS."]}),"\n",(0,t.jsxs)(n.li,{children:["Remove ",(0,t.jsx)(n.code,{children:"retry_timeout"})," and other HTTP-only keys from ",(0,t.jsx)(n.code,{children:"client.conf.string"}),".\nKeep your credentials and data mapping settings."]}),"\n",(0,t.jsxs)(n.li,{children:["Enable ",(0,t.jsx)(n.a,{href:"#exactly-once-delivery",children:"deduplication"})," if duplicate events are not\nacceptable, and review ",(0,t.jsx)(n.a,{href:"#outages-and-reconnects",children:"outage recovery"}),"."]}),"\n"]}),"\n",(0,t.jsx)(n.h3,{id:"configuration-reference",children:"Configuration reference"}),"\n",(0,t.jsxs)(n.p,{children:["Set the QuestDB address and credentials in ",(0,t.jsx)(n.code,{children:"client.conf.string"}),". Add data\nmapping and delivery options as separate connector properties."]}),"\n",(0,t.jsx)(n.h4,{id:"connector-options",children:"Connector options"}),"\n",(0,t.jsxs)(n.table,{children:[(0,t.jsx)(n.thead,{children:(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.th,{children:"Name"}),(0,t.jsx)(n.th,{children:"Type"}),(0,t.jsx)(n.th,{children:"Example"}),(0,t.jsx)(n.th,{children:"Default"}),(0,t.jsx)(n.th,{children:"Description"})]})}),(0,t.jsxs)(n.tbody,{children:[(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"client.conf.string"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"string"})}),(0,t.jsx)(n.td,{children:"ws::addr=localhost:9000;"}),(0,t.jsx)(n.td,{children:"N/A"}),(0,t.jsx)(n.td,{children:"Client configuration string"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"topics"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"string"})}),(0,t.jsx)(n.td,{children:"orders,audit"}),(0,t.jsx)(n.td,{children:"N/A"}),(0,t.jsx)(n.td,{children:"Kafka topics to read from"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"table"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"string"})}),(0,t.jsx)(n.td,{children:"my_table"}),(0,t.jsx)(n.td,{children:"Topic name"}),(0,t.jsx)(n.td,{children:"Target table in QuestDB"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"key.converter"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"string"})}),(0,t.jsx)(n.td,{children:(0,t.jsx)("sub",{children:"org.apache.kafka.connect.storage.StringConverter"})}),(0,t.jsx)(n.td,{children:"N/A"}),(0,t.jsx)(n.td,{children:"Converter for Kafka keys"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"value.converter"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"string"})}),(0,t.jsx)(n.td,{children:(0,t.jsx)("sub",{children:"org.apache.kafka.connect.json.JsonConverter"})}),(0,t.jsx)(n.td,{children:"N/A"}),(0,t.jsx)(n.td,{children:"Converter for Kafka values"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"value.format"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"string"})}),(0,t.jsx)(n.td,{children:"json"}),(0,t.jsx)(n.td,{children:"connect"}),(0,t.jsxs)(n.td,{children:["Payload format: ",(0,t.jsx)(n.code,{children:"connect"}),", ",(0,t.jsx)(n.code,{children:"json"}),", or ",(0,t.jsx)(n.code,{children:"json_envelope"}),". See ",(0,t.jsx)(n.a,{href:"#raw-json-fast-path",children:"Raw JSON fast path"})]})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"include.key"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"boolean"})}),(0,t.jsx)(n.td,{children:"false"}),(0,t.jsx)(n.td,{children:"true"}
1),(0,t.jsx)(n.td,{children:"Include message key in target table"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"key.prefix"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"string"})}),(0,t.jsx)(n.td,{children:"from_key"}),(0,t.jsx)(n.td,{children:"key"}),(0,t.jsx)(n.td,{children:"Prefix for key fields"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"value.prefix"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"string"})}),(0,t.jsx)(n.td,{children:"from_value"}),(0,t.jsx)(n.td,{children:"N/A"}),(0,t.jsx)(n.td,{children:"Prefix for value fields"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"symbols"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"string"})}),(0,t.jsx)(n.td,{children:"instrument,stock"}),(0,t.jsx)(n.td,{children:"N/A"}),(0,t.jsxs)(n.td,{children:["Columns to create as ",(0,t.jsx)(n.a,{href:"/docs/concepts/symbol/",children:"symbol"})," type"]})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"doubles"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"string"})}),(0,t.jsx)(n.td,{children:"volume,price"}),(0,t.jsx)(n.td,{children:"N/A"}),(0,t.jsx)(n.td,{children:"Columns to always send as double type"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"timestamp.field.name"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"string"})}),(0,t.jsx)(n.td,{children:"pickup_time"}),(0,t.jsx)(n.td,{children:"N/A"}),(0,t.jsxs)(n.td,{children:["Designated timestamp field. Use comma-separated names for ",(0,t.jsx)(n.a,{href:"#composed-timestamps",children:"composed timestamps"})]})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"timestamp.units"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"string"})}),(0,t.jsx)(n.td,{children:"micros"}),(0,t.jsx)(n.td,{children:"auto"}),(0,t.jsxs)(n.td,{children:["Timestamp field units: ",(0,t.jsx)(n.code,{children:"nanos"}),", ",(0,t.jsx)(n.code,{children:"micros"}),", ",(0,t.jsx)(n.code,{children:"millis"}),", ",(0,t.jsx)(n.code,{children:"seconds"}),", ",(0,t.jsx)(n.code,{children:"auto"})]})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"timestamp.kafka.native"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"boolean"})}),(0,t.jsx)(n.td,{children:"true"}),(0,t.jsx)(n.td,{children:"false"}),(0,t.jsx)(n.td,{children:"Use Kafka message timestamps as designated timestamps"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"timestamp.string.fields"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"string"})}),(0,t.jsx)(n.td,{children:"creation_time"}),(0,t.jsx)(n.td,{children:"N/A"}),(0,t.jsx)(n.td,{children:"String fields containing textual timestamps"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"timestamp.string.format"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"string"})}),(0,t.jsxs)(n.td,{children:["yyyy-MM-dd HH:mm",":ss",".SSSUUU z"]}),(0,t.jsx)(n.td,{children:(0,t.jsxs)("sub",{children:["yyyy-MM-ddTHH:mm",":ss",".SSSUUUZ"]})}),(0,t.jsx)(n.td,{children:"Format for parsing string timestamps"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"skip.unsupported.types"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"boolean"})}),(0,t.jsx)(n.td,{children:"false"}),(0,t.jsx)(n.td,{children:"false"}),(0,t.jsx)(n.td,{children:"Skip unsupported types instead of failing"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"allowed.lag"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"int"})}),(0,t.jsx)(n.td,{children:"250"}),(0,t.jsx)(n.td,{children:"1000"}),(0,t.jsx)(n.td,{children:"Maximum wait in milliseconds for new records before sending a partial batch when idle"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"retry.backoff.ms"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"long"})}),(0,t.jsx)(n.td,{children:"5000"}),(0,t.jsx)(n.td,{children:"3000"}),(0,t.jsx)(n.td,{children:"Milliseconds to wait before reconnecting when QuestDB is unreachable. Not used by the HTTP transport"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"dlq.send.batch.on.error"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"boolean"})}),(0,t.jsx)(n.td,{children:"true"}),(0,t.jsx)(n.td,{children:"false"}),(0,t.jsxs)(n.td,{children:["Send a whole rejected batch to the dead letter queue, including any valid records in it. See ",(0,t.jsx)(n.a,{href:"#dead-letter-queue",children:"Dead letter queue"})]})]})]})]}),"\n",(0,t.jsxs)(n.p,{children:["The connector uses Kafka Connect converters for deserialization and works with\nany format they support, including JSON, Avro, and Protobuf. When using Schema\nRegistry, configure the appropriate converter (e.g.,\n",(0,t.jsx)(n.code,{children:"io.confluent.connect.avro.AvroConverter"}),")."]}),"\n",(0,t.jsx)(n.h4,{id:"qwp-delivery-options",children:"QWP delivery options"}),"\n",(0,t.jsxs)(n.p,{children:["These options apply only to ",(0,t.jsx)(n.code,{children:"ws"})," and ",(0,t.jsx)(n.code,{children:"wss"}),". Start with the defaults; see\n",(0,t.jsx)(n.a,{href:"#outages-and-reconnects",children:"outage recovery"})," and\n",(0,t.jsx)(n.a,{href:"#performance-tuning",children:"performance tuning"})," before changing them."]}),"\n",(0,t.jsxs)(n.table,{children:[(0,t.jsx)(n.thead,{children:(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.th,{children:"Name"}),(0,t.jsx)(n.th,{children:"Type"}),(0,t.jsx)(n.th,{children:"Default"}),(0,t.jsx)(n.th,{children:"Description"})]})}),(0,t.jsxs)(n.tbody,{children:[(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"qwp.max.inflight.rows"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"int"})}),(0,t.jsx)(n.td,{children:"150000"}),(0,t.jsx)(n.td,{children:"Soft limit on buffered or sent rows awaiting confirmation. Consumption pauses above this limit; the current Kafka poll batch can exceed it"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"qwp.progress.timeout.ms"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"long"})}),(0,t.jsx)(n.td,{children:"300000"}),(0,t.jsx)(n.td,{children:"Milliseconds without delivery progress before the task fails while rows are pending. New acknowledgements reset the timer"})]})]})]}),"\n",(0,t.jsxs)(s,{children:[(0,t.jsx)("summary",{children:"Advanced delivery options"}),(0,t.jsxs)(n.table,{children:[(0,t.jsx)(n.thead,{children:(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.th,{children:"Name"}),(0,t.jsx)(n.th,{children:"Type"}),(0,t.jsx)(n.th,{children:"Default"}),(0,t.jsx)(n.th,{children:"Description"})]})}),(0,t.jsxs)(n.tbody,{children:[(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"qwp.commit.ack.timeout.ms"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"long"})}),(0,t.jsx)(n.td,{children:"500"}),(0,t.jsx)(n.td,{children:"How long an offset commit waits for delivery confirmation, in milliseconds. On timeout, unconfirmed offsets remain uncommitted; this alone does not trigger redelivery"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"qwp.dlq.terminal.categories"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"list"})}),(0,t.jsx)(n.td,{children:"SCHEMA_MISMATCH"}),(0,t.jsx)(n.td,{children:"Server errors eligible for the DLQ. Keep the default to avoid treating infrastru
1cture failures as bad records"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:"qwp.quarantine.ack.timeout.ms"}),(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"long"})}),(0,t.jsx)(n.td,{children:"1000"}),(0,t.jsx)(n.td,{children:"How long each batch waits for delivery confirmation while isolating a rejected record, in milliseconds"})]})]})]}),(0,t.jsxs)(n.p,{children:["Pre-release builds of the QWP transport used the names ",(0,t.jsx)(n.code,{children:"max.inflight.rows"})," and\n",(0,t.jsx)(n.code,{children:"progress.timeout.ms"}),". Rename them to\n",(0,t.jsx)(n.code,{children:"qwp.max.inflight.rows"})," and ",(0,t.jsx)(n.code,{children:"qwp.progress.timeout.ms"}),"."]})]}),"\n",(0,t.jsx)(n.h4,{id:"client-configuration-string",children:"Client configuration string"}),"\n",(0,t.jsxs)(n.p,{children:["The ",(0,t.jsx)(n.code,{children:"client.conf.string"})," option configures how the connector communicates with\nQuestDB. You can also set this via the ",(0,t.jsx)(n.code,{children:"QDB_CLIENT_CONF"})," environment variable."]}),"\n",(0,t.jsx)(n.p,{children:"Format:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{children:"<protocol>::<key>=<value>;<key>=<value>;...;\n"})}),"\n",(0,t.jsx)(n.p,{children:"Note the trailing semicolon."}),"\n",(0,t.jsx)(n.p,{children:(0,t.jsx)(n.strong,{children:"Supported protocols:"})}),"\n",(0,t.jsxs)(n.table,{children:[(0,t.jsx)(n.thead,{children:(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.th,{children:"Protocol"}),(0,t.jsx)(n.th,{children:"Transport"}),(0,t.jsx)(n.th,{children:"Notes"})]})}),(0,t.jsxs)(n.tbody,{children:[(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"ws"})}),(0,t.jsx)(n.td,{children:"QWP over WebSocket"}),(0,t.jsx)(n.td,{children:"Recommended. Acknowledged delivery, automatic reconnects"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"wss"})}),(0,t.jsx)(n.td,{children:"QWP over WebSocket with TLS"}),(0,t.jsx)(n.td,{children:"Requires QuestDB Enterprise, or a TLS-terminating proxy in front of QuestDB open source"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsxs)(n.td,{children:[(0,t.jsx)(n.code,{children:"http"}),", ",(0,t.jsx)(n.code,{children:"https"})]}),(0,t.jsx)(n.td,{children:"ILP over HTTP"}),(0,t.jsxs)(n.td,{children:["Legacy. See ",(0,t.jsx)(n.a,{href:"#legacy-ilp-transports",children:"Legacy ILP transports"})]})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsxs)(n.td,{children:[(0,t.jsx)(n.code,{children:"tcp"}),", ",(0,t.jsx)(n.code,{children:"tcps"})]}),(0,t.jsx)(n.td,{children:"ILP over TCP"}),(0,t.jsx)(n.td,{children:"Not recommended. Offers no delivery guarantees"})]})]})]}),"\n",(0,t.jsx)(n.p,{children:(0,t.jsx)(n.strong,{children:"Required keys:"})}),"\n",(0,t.jsxs)(n.ul,{children:["\n",(0,t.jsxs)(n.li,{children:[(0,t.jsx)(n.code,{children:"addr"})," - QuestDB hostname and port (port defaults to 9000)"]}),"\n"]}),"\n",(0,t.jsx)(n.p,{children:"Examples:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"# Minimal configuration\nclient.conf.string=ws::addr=localhost:9000;\n\n# Basic authentication with the password from an environment variable\nclient.conf.string=ws::addr=questdb.example.com:9000;username=admin;password=${QUESTDB_PASSWORD};\n\n# TLS with a bearer token (QuestDB Enterprise)\nclient.conf.string=wss::addr=questdb.example.com:9000;token=${QUESTDB_TOKEN};\n\n# Multi-host failover (QuestDB Enterprise)\nclient.conf.string=wss::addr=node-a:9000,node-b:9000;token=${QUESTDB_TOKEN};\n"})}),"\n",(0,t.jsxs)(n.p,{children:["See the ",(0,t.jsx)(n.a,{href:"/docs/connect/clients/connect-string/",children:"connect string reference"})," for\nall available client keys."]}),"\n",(0,t.jsx)(n.h5,{id:"batching-and-buffer-options",children:"Batching and buffer options"}),"\n",(0,t.jsxs)(n.p,{children:["These client settings apply to ",(0,t.jsx)(n.code,{children:"ws"})," and ",(0,t.jsx)(n.code,{children:"wss"}),". For examples, see\n",(0,t.jsx)(n.a,{href:"#performance-tuning",children:"performance tuning"}),"."]}),"\n",(0,t.jsxs)(n.table,{children:[(0,t.jsx)(n.thead,{children:(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.th,{children:"Key"}),(0,t.jsx)(n.th,{children:"Behaviour in the connector"})]})}),(0,t.jsxs)(n.tbody,{children:[(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"auto_flush_rows"})}),(0,t.jsxs)(n.td,{children:["Send a batch at this many rows. Default: ",(0,t.jsx)(n.code,{children:"75000"}),". Cannot be ",(0,t.jsx)(n.code,{children:"off"})]})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"auto_flush_interval"})}),(0,t.jsxs)(n.td,{children:["Interval for sending pending rows, in milliseconds. Default: ",(0,t.jsx)(n.code,{children:"1000"}),". Cannot be ",(0,t.jsx)(n.code,{children:"off"})]})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"sf_max_total_bytes"})}),(0,t.jsxs)(n.td,{children:["Cap on the memory buffer of encoded, unacknowledged rows. Default: ",(0,t.jsx)(n.code,{children:"128m"})]})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"sf_append_deadline_millis"})}
1),(0,t.jsxs)(n.td,{children:["How long sending can wait for buffer space, in milliseconds. Default: ",(0,t.jsx)(n.code,{children:"30000"}),". Must be lower than the consumer's ",(0,t.jsx)(n.code,{children:"max.poll.interval.ms"})]})]})]})]}),"\n",(0,t.jsxs)(s,{children:[(0,t.jsx)("summary",{children:"Client settings with connector-specific behavior"}),(0,t.jsxs)(n.ul,{children:["\n",(0,t.jsxs)(n.li,{children:["Leave ",(0,t.jsx)(n.code,{children:"auto_flush_bytes"})," enabled so batches fit the server's size limit."]}),"\n",(0,t.jsxs)(n.li,{children:["Omit ",(0,t.jsx)(n.code,{children:"sf_dir"})," and ",(0,t.jsx)(n.code,{children:"sf_durability"}),": disk buffering is not supported. Recovery\nrelies on ",(0,t.jsx)(n.a,{href:"#outages-and-reconnects",children:"Kafka retention"}),"."]}),"\n",(0,t.jsxs)(n.li,{children:["Omit ",(0,t.jsx)(n.code,{children:"initial_connect_retry"})," or set it to ",(0,t.jsx)(n.code,{children:"off"}),". The connector handles\nstartup retries using ",(0,t.jsx)(n.code,{children:"retry.backoff.ms"}),"; ",(0,t.jsx)(n.code,{children:"reconnect_*"})," settings control\nretries after a connection drops."]}),"\n",(0,t.jsxs)(n.li,{children:["Leave ",(0,t.jsx)(n.code,{children:"close_flush_timeout_millis"})," at its default of ",(0,t.jsx)(n.code,{children:"0"}),". Increasing it\ndelays shutdown without allowing more offsets to be committed."]}),"\n"]})]}),"\n",(0,t.jsx)(n.h5,{id:"environment-variable-expansion",children:"Environment variable expansion"}),"\n",(0,t.jsxs)(n.p,{children:["The ",(0,t.jsx)(n.code,{children:"client.conf.string"})," supports ",(0,t.jsx)(n.code,{children:"${VAR}"})," syntax for environment variable\nexpansion, useful for injecting secrets in Kubernetes environments:"]}),"\n",(0,t.jsxs)(n.table,{children:[(0,t.jsx)(n.thead,{children:(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.th,{children:"Pattern"}),(0,t.jsx)(n.th,{children:"Result"})]})}),(0,t.jsxs)(n.tbody,{children:[(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"${VAR}"})}),(0,t.jsx)(n.td,{children:"Replaced with environment variable value"})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"$$"})}),(0,t.jsxs)(n.td,{children:["Escaped to literal ",(0,t.jsx)(n.code,{children:"$"})]})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"$${VAR}"})}),(0,t.jsxs)(n.td,{children:["Escaped to literal ",(0,t.jsx)(n.code,{children:"${VAR}"})," (not expanded)"]})]}),(0,t.jsxs)(n.tr,{children:[(0,t.jsx)(n.td,{children:(0,t.jsx)(n.code,{children:"$VAR"})}),(0,t.jsx)(n.td,{children:"Not expanded (braces required)"})]})]})]}),"\n",(0,t.jsx)(n.p,{children:"The connector fails to start if:"}),"\n",(0,t.jsxs)(n.ul,{children:["\n",(0,t.jsx)(n.li,{children:"A referenced environment variable is not defined"}),"\n",(0,t.jsxs)(n.li,{children:["A variable reference is malformed (e.g., unclosed braces ",(0,t.jsx)(n.code,{children:"${VAR"}),")"]}),"\n",(0,t.jsxs)(n.li,{children:["A variable name is empty (",(0,t.jsx)(n.code,{children:"${}"}),") or invalid (must start with letter or\nunderscore, followed by letters, digits, or underscores)"]}),"\n"]}),"\n",(0,t.jsx)(n.admonition,{type:"warning",children:(0,t.jsxs)(n.p,{children:["Environment variable values containing semicolons (",(0,t.jsx)(n.code,{children:";"}),") will break the\nconfiguration string parsing."]})}),"\n",(0,t.jsx)(n.h3,{id:"sample-projects",children:"Sample projects"}),"\n",(0,t.jsx)(n.p,{children:"Additional examples are available on GitHub:"}),"\n",(0,t.jsxs)(n.ul,{children:["\n",(0,t.jsx)(n.li,{children:(0,t.jsx)(n.a,{href:"https://github.com/questdb/kafka-questdb-connector/tree/main/kafka-questdb-connector-samples",children:"Sample projects"})}),"\n",(0,t.jsx)(n.li,{children:(0,t.jsx)(n.a,{href:"https://github.com/questdb/kafka-questdb-connector/tree/main/kafka-questdb-connector-samples/stocks",children:"Debezium CDC integration"})}),"\n"]}),"\n",(0,t.jsx)(n.h2,{id:"stream-processing",children:"Stream processing"}),"\n",(0,t.jsxs)(n.p,{children:[(0,t.jsx)(n.a,{href:"/glossary/stream-processing/",children:"Stream processing"})," engines like\n",(0,t.jsx)(n.a,{href:"https://flink.apache.org/",children:"Apache Flink"})," provide rich APIs for data\ntransformation, enrichment, and filtering with built-in fault tolerance."]}),"\n",(0,t.jsxs)(n.p,{children:["QuestDB offers a ",(0,t.jsx)(n.a,{href:"/docs/connect/message-brokers/flink/",children:"Flink connector"})," for\nusers who need complex transformations while ingesting from Kafka."]}),"\n",(0,t.jsx)(n.p,{children:(0,t.jsx)(n.strong,{children:"Use stream processing when you need:"})}),"\n",(0,t.jsxs)(n.ul,{children:["\n",(0,t.jsx)(n.li,{children:"Complex stateful transformations"}),"\n",(0,t.jsx)(n.li,{children:"Joining multiple data streams"}),"\n",(0,t.jsx)(n.li,{children:"Windowed aggregations before writing to QuestDB"}),"\n"]}),"\n",(0,t.jsx)(n.h2,{id:"custom-program",children:"Custom program"}),"\n",(0,t.jsx)(n.p,{children:"Writing a dedicated program to read from Kafka and write to QuestDB offers\nmaximum flexibility for arbitrary transformations and filtering."}),"\n",(0,t.jsx)(n.p,{children:(0,t.jsx)(n.strong,{children:"Trade-offs:"})}),"\n",(0,t.jsxs)(n.ul,{children:["\n",(0,t.jsx)(n.li,{children:"Full control over serialization, error handling, and batching"}
1),"\n",(0,t.jsx)(n.li,{children:"Highest implementation complexity"}),"\n",(0,t.jsx)(n.li,{children:"Must handle Kafka consumer groups, offset management, and retries"}),"\n"]}),"\n",(0,t.jsx)(n.p,{children:"This approach is only recommended for advanced use cases where the Kafka\nconnector or stream processing cannot meet your requirements."}),"\n",(0,t.jsx)(n.h2,{id:"faq",children:"FAQ"}),"\n",(0,t.jsxs)(s,{children:[(0,t.jsx)("summary",{children:"Does the connector work with Schema Registry?"}),(0,t.jsxs)(n.p,{children:["Yes. The connector relies on Kafka Connect converters for deserialization.\nConfigure converters using ",(0,t.jsx)(n.code,{children:"key.converter"})," and ",(0,t.jsx)(n.code,{children:"value.converter"})," options.\nIt works with Avro, JSON Schema, and other formats supported by Schema Registry."]})]}),"\n",(0,t.jsxs)(s,{children:[(0,t.jsx)("summary",{children:"Does the connector work with Debezium?"}),(0,t.jsxs)(n.p,{children:["Yes. QuestDB works well with ",(0,t.jsx)(n.a,{href:"https://debezium.io/",children:"Debezium"})," for\n",(0,t.jsx)(n.a,{href:"/glossary/change-data-capture/",children:"change data capture"}),". Since QuestDB is\nappend-only, updates become new rows preserving history."]}),(0,t.jsxs)(n.p,{children:["Use Debezium's ",(0,t.jsx)(n.code,{children:"ExtractNewRecordState"})," transformation to extract the new record\nstate. DELETE events are dropped by default."]}),(0,t.jsxs)(n.p,{children:["See the ",(0,t.jsx)(n.a,{href:"https://github.com/questdb/kafka-questdb-connector/tree/main/kafka-questdb-connector-samples/stocks",children:"Debezium sample project"}),"\nand the blog post ",(0,t.jsx)(n.a,{href:"/blog/2023/01/03/change-data-capture-with-questdb-and-debezium/",children:"Change Data Capture with QuestDB and Debezium"}),"."]}),(0,t.jsxs)(n.p,{children:[(0,t.jsx)(n.strong,{children:"Typical pattern:"})," Use a relational database for current state and QuestDB\nfor change history. For example, PostgreSQL holds current stock prices while\nQuestDB stores the complete price history for analytics."]})]}),"\n",(0,t.jsxs)(s,{children:[(0,t.jsx)("summary",{children:"How do I select which fields to include?"}),(0,t.jsxs)(n.p,{children:["Use Kafka Connect's ",(0,t.jsx)(n.code,{children:"ReplaceField"})," transformation:"]}),(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-json",children:'{\n  "transforms": "removeFields",\n  "transforms.removeFields.type": "org.apache.kafka.connect.transforms.ReplaceField$Value",\n  "transforms.removeFields.blacklist": "address,internal_id"\n}\n'})}),(0,t.jsxs)(n.p,{children:["See ",(0,t.jsx)(n.a,{href:"https://docs.confluent.io/platform/current/connect/transforms/replacefield.html",children:"ReplaceField documentation"}),"."]})]}),"\n",(0,t.jsxs)(s,{children:[(0,t.jsx)("summary",{children:"I'm getting a JsonConverter schema error"}),(0,t.jsx)(n.p,{children:"If you see:"}),(0,t.jsxs)(n.blockquote,{children:["\n",(0,t.jsx)(n.p,{children:"JsonConverter with schemas.enable requires 'schema' and 'payload' fields"}),"\n"]}),(0,t.jsx)(n.p,{children:"Your JSON data doesn't include a schema. Add to your configuration:"}),(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"value.converter.schemas.enable=false\n"})}),(0,t.jsx)(n.p,{children:"Or for keys:"}),(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-properties",children:"key.converter.schemas.enable=false\n"})})]}),"\n",(0,t.jsxs)(s,{children:[(0,t.jsx)("summary",{children:'The task fails with "QWP acknowledgements did not advance"'}),(0,t.jsxs)(n.p,{children:["QuestDB has not confirmed further delivery for ",(0,t.jsx)(n.code,{children:"qwp.progress.timeout.ms"}),"\n(default 5 minutes) while rows were pending. Check that QuestDB is available\nand reachable from the Kafka Connect worker, then restart the task. See\n",(0,t.jsx)(n.a,{href:"#outages-and-reconnects",children:"outage recovery"})," for timeout and retention settings."]})]}),"\n",(0,t.jsx)(n.h2,{id:"see-also",children:"See also"}),"\n",(0,t.jsxs)(n.ul,{children:["\n",(0,t.jsx)(n.li,{children:(0,t.jsx)(n.a,{href:"/docs/concepts/delivery-semantics/",children:"Delivery semantics"})}),"\n",(0,t.jsx)(n.li,{children:(0,t.jsx)(n.a,{href:"/docs/connect/clients/connect-string/",children:"Connect string reference"})}),"\n",(0,t.jsx)(n.li,{children:(0,t.jsx)(n.a,{href:"/blog/2023/01/03/change-data-capture-with-questdb-and-debezium/",children:"Change Data Capture with QuestDB and Debezium"})}),"\n",(0,t.jsx)(n.li,{children:(0,t.jsx)(n.a,{href:"/blog/realt
1ime-crypto-tracker-with-questdb-kafka-connector/",children:"Realtime crypto tracker with QuestDB Kafka Connector"})}),"\n"]})]})}function h(e={}){let{wrapper:n}={...(0,i.a)(),...e.components};return n?(0,t.jsx)(n,{...e,children:(0,t.jsx)(l,{...e})}):l(e)}},50065:function(e,n,s){s.d(n,{Z:()=>c,a:()=>d});var r=s(67294);let t={},i=r.createContext(t);function d(e){let n=r.useContext(i);return r.useMemo(function(){return"function"==typeof e?e(n):{...n,...e}},[n,e])}function c(e){let n;return n=e.disableParentContext?"function"==typeof e.components?e.components(t):e.components||t:d(e.components),r.createElement(i.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.