{"_id":"@danhawkins/avro-kafkajs","_rev":"4-cb2d2d5b8604fa471d1b56f96591d5e9","name":"@danhawkins/avro-kafkajs","dist-tags":{"latest":"0.8.2"},"versions":{"0.8.0":{"name":"@danhawkins/avro-kafkajs","version":"0.8.0","main":"dist/index.js","types":"dist/index.d.ts","description":"A wrapper around Kafkajs to transparently use Schema Registry for producing and consuming messages with avro schemas.","author":{"name":"Ivan Kerin","email":"ikerin@gmail.com"},"repository":{"type":"git","url":"git@github.com:ovotech/castle.git"},"homepage":"https://github.com/ovotech/castle/tree/main/packages/avro-kafkajs#readme","license":"Apache-2.0","devDependencies":{"@ovotech/build-docs":"^0.1.0","@types/jest":"^26.0.20","@types/long":"^4.0.1","@types/node":"^14.14.28","@types/uuid":"^8.3.0","@typescript-eslint/eslint-plugin":"^4.15.1","@typescript-eslint/parser":"^4.15.1","axios":"^0.21.0","eslint":"^7.20.0","eslint-config-prettier":"^7.2.0","jest":"^26.6.3","kafkajs":"^1.15.0","prettier":"^2.2.1","stream-mock":"^2.0.5","ts-jest":"^26.5.1","ts-node":"^9.1.1","ts-retry-promise":"^0.6.0","typescript":"^4.1.2","uuid":"^8.3.1"},"scripts":{"build:docs":"build-docs README.md","build":"tsc --declaration","test":"jest test --runInBand","lint:prettier":"prettier --list-different {src,test}/**/*.ts","lint:eslint":"eslint '{src,test}/**/*.ts'","lint":"yarn lint:prettier && yarn lint:eslint"},"jest":{"preset":"../../jest.json"},"peerDependencies":{"kafkajs":"^1.15.0"},"dependencies":{"@ovotech/schema-registry-api":"^1.1.1","avsc":"^5.5.6","long":"^4.0.0"},"licenseText":"Copyright 2019-2020 OVO Energy\n\nLicensed under the Apache License, Version 2.0 (the \"License\");\nyou may not use this file except in compliance with the License.\nYou may obtain a copy of the License at\n\nhttp://www.apache.org/licenses/LICENSE-2.0\n\nUnless required by applicable law or agreed to in writing, software\ndistributed under the License is distributed on an \"AS IS\" BASIS,\nWITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.\nSee the License for the specific language governing permissions and\nlimitations under the License.\n","_id":"@danhawkins/avro-kafkajs@0.8.0","dist":{"shasum":"adb8761c2c9b7f404d9209ae2db33f78a0c3481e","integrity":"sha512-uttsTzEqN82LgWuXrDpeHswn/7SuQe9p+GHhBfk9FWaucY5v8SRH8ke1aAdJFcD4pMdot0D6lWsygP/zXn3B1Q==","tarball":"https://registry.npmjs.org/@danhawkins/avro-kafkajs/-/avro-kafkajs-0.8.0.tgz","fileCount":23,"unpackedSize":51587,"signatures":[{"keyid":"SHA256:jl3bwswu80PjjokCgh0o2w5c2U4LhQAE57gj9cz1kzA","sig":"MEQCIFAGJLfMmitRPOSUt5lEaUTBIT8YYbwB4kMBvAD6J0xlAiBaIXhzyYIn4L3AQOLkCBfl7i7PJCIg8BYwU37/GVOgWg=="}]},"_npmUser":{"name":"danhawkins","email":"danny@quiqup.com"},"directories":{},"maintainers":[{"name":"danhawkins","email":"danny@quiqup.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages","tmp":"tmp/avro-kafkajs_0.8.0_1632936465231_0.6213539960306795"},"_hasShrinkwrap":false},"0.8.1":{"name":"@danhawkins/avro-kafkajs","version":"0.8.1","main":"dist/index.js","types":"dist/index.d.ts","description":"A wrapper around Kafkajs to transparently use Schema Registry for producing and consuming messages with avro schemas.","author":{"name":"Ivan Kerin","email":"ikerin@gmail.com"},"repository":{"type":"git","url":"git@github.com:ovotech/castle.git"},"homepage":"https://github.com/ovotech/castle/tree/main/packages/avro-kafkajs#readme","license":"Apache-2.0","devDependencies":{"@ovotech/build-docs":"^0.1.0","@types/jest":"^26.0.20","@types/long":"^4.0.1","@types/node":"^14.14.28","@types/uuid":"^8.3.0","@typescript-eslint/eslint-plugin":"^4.15.1","@typescript-eslint/parser":"^4.15.1","axios":"^0.21.0","eslint":"^7.20.0","eslint-config-prettier":"^7.2.0","jest":"^26.6.3","kafkajs":"^1.15.0","prettier":"^2.2.1","stream-mock":"^2.0.5","ts-jest":"^26.5.1","ts-node":"^9.1.1","ts-retry-promise":"^0.6.0","typescript":"^4.1.2","uuid":"^8.3.1"},"scripts":{"build:docs":"build-docs README.md","build":"tsc --declaration","test":"jest test --runInBand","lint:prettier":"prettier --list-different {src,test}/**/*.ts","lint:eslint":"eslint '{src,test}/**/*.ts'","lint":"yarn lint:prettier && yarn lint:eslint"},"jest":{"preset":"../../jest.json"},"peerDependencies":{"kafkajs":"^1.15.0"},"dependencies":{"@ovotech/schema-registry-api":"^1.1.1","avsc":"^5.5.6","long":"^4.0.0"},"licenseText":"Copyright 2019-2020 OVO Energy\n\nLicensed under the Apache License, Version 2.0 (the \"License\");\nyou may not use this file except in compliance with the License.\nYou may obtain a copy of the License at\n\nhttp://www.apache.org/licenses/LICENSE-2.0\n\nUnless required by applicable law or agreed to in writing, software\ndistributed under the License is distributed on an \"AS IS\" BASIS,\nWITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.\nSee the License for the specific language governing permissions and\nlimitations under the License.\n","_id":"@danhawkins/avro-kafkajs@0.8.1","dist":{"shasum":"7f34010ce7d8f23ffc9bcf2b71caff956fefd84c","integrity":"sha512-dYMNitfX77nHvFbx8yP3PlEmPtnKKcXWbdPdRqWmbow/9uxDhKhlDVzaOVPQP7UaGf3kSM8pWDK0blP0ggBvQA==","tarball":"https://registry.npmjs.org/@danhawkins/avro-kafkajs/-/avro-kafkajs-0.8.1.tgz","fileCount":23,"unpackedSize":51611,"signatures":[{"keyid":"SHA256:jl3bwswu80PjjokCgh0o2w5c2U4LhQAE57gj9cz1kzA","sig":"MEUCIDECCt0iTd1GK6L/B4cp0wgCCwHTNatffYNVbe6KV2GDAiEA/PbUXQy5C/GPKWYma/1Ck1Q5HSyZdZgcdVEbfAOnVwA="}]},"_npmUser":{"name":"danhawkins","email":"danny@quiqup.com"},"directories":{},"maintainers":[{"name":"danhawkins","email":"danny@quiqup.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages","tmp":"tmp/avro-kafkajs_0.8.1_1632936700065_0.5403450942235901"},"_hasShrinkwrap":false},"0.8.2":{"name":"@danhawkins/avro-kafkajs","version":"0.8.2","main":"dist/index.js","types":"dist/index.d.ts","description":"A wrapper around Kafkajs to transparently use Schema Registry for producing and consuming messages with avro schemas.","author":{"name":"Ivan Kerin","email":"ikerin@gmail.com"},"repository":{"type":"git","url":"git@github.com:ovotech/castle.git"},"homepage":"https://github.com/ovotech/castle/tree/main/packages/avro-kafkajs#readme","license":"Apache-2.0","devDependencies":{"@ovotech/build-docs":"^0.1.0","@types/jest":"^26.0.20","@types/long":"^4.0.1","@types/node":"^14.14.28","@types/uuid":"^8.3.0","@typescript-eslint/eslint-plugin":"^4.15.1","@typescript-eslint/parser":"^4.15.1","axios":"^0.21.0","eslint":"^7.20.0","eslint-config-prettier":"^7.2.0","jest":"^26.6.3","kafkajs":"^1.15.0","prettier":"^2.2.1","stream-mock":"^2.0.5","ts-jest":"^26.5.1","ts-node":"^9.1.1","ts-retry-promise":"^0.6.0","typescript":"^4.1.2","uuid":"^8.3.1"},"scripts":{"build:docs":"build-docs README.md","build":"tsc --declaration","test":"jest test --runInBand","lint:prettier":"prettier --list-different {src,test}/**/*.ts","lint:eslint":"eslint '{src,test}/**/*.ts'","lint":"yarn lint:prettier && yarn lint:eslint"},"jest":{"preset":"../../jest.json"},"peerDependencies":{"kafkajs":"^1.15.0"},"dependencies":{"@ovotech/schema-registry-api":"^1.1.1","avsc":"^5.5.6","long":"^4.0.0"},"licenseText":"Copyright 2019-2020 OVO Energy\n\nLicensed under the Apache License, Version 2.0 (the \"License\");\nyou may not use this file except in compliance with the License.\nYou may obtain a copy of the License at\n\nhttp://www.apache.org/licenses/LICENSE-2.0\n\nUnless required by applicable law or agreed to in writing, software\ndistributed under the License is distributed on an \"AS IS\" BASIS,\nWITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.\nSee the License for the specific language governing permissions and\nlimitations under the License.\n","_id":"@danhawkins/avro-kafkajs@0.8.2","dist":{"shasum":"74fd12b9d7b71b59677d186f9a441a498bc1852b","integrity":"sha512-xFJOeghDts1fi8ioV2RxMOkEebaBJ+SX5QTCV9E1pEY+OLcbxJWg58xSDTPkhavs04ZHXXE2woXE1erpGFNSJg==","tarball":"https://registry.npmjs.org/@danhawkins/avro-kafkajs/-/avro-kafkajs-0.8.2.tgz","fileCount":23,"unpackedSize":51611,"npm-signature":"-----BEGIN PGP SIGNATURE-----\r\nVersion: OpenPGP.js v3.0.13\r\nComment: https://openpgpjs.org\r\n\r\nwsFcBAEBCAAQBQJhzewVCRA9TVsSAnZWagAAjpoP/2oZCyeykSWgrpVQ3Thg\nPbkwmq5QkGoFKWAcMwrAoCNwdDBZcn4XNUX0Bq50ZQFbixcVk0GztFZ/iLi4\nthGm48EZ4ixQ3IWfdjlk+17h3kLTQTU5fQOsF0EmUl+QeuW2Ak99oLhlzxSN\nQ4Spyvy3wL7PajLF2zp8OgGTLK0/N+hyuLn0frb5ag4zRB9X20Q2JYbOEF9Z\nsFvaiDLD0tN2I4udqhyoovKQNN7oeckv14/jeSIyIQWiywzNWWdlYxbS+jY0\n2G1qBP6VyZ3qgC7+EmXoe6k6p8jUAN08T1yBz7i/sINsZDJ4lz7rvHgLpTE3\nDZHqfikSeU1rQq8UNU0Xp8Aibw/Nd9IRfcX70s80yjImATrusPNwKrzNwZnS\ndz1Fd05fQdXw7FuhdvC4j08xT4C+b7YSWAUBdfuvTxPrscwj3z7f2tND+QkI\nEcFpg8B///8mIYTWWBqftv+gu22khNv/CYEicosq204GR3whiS6tlKsdZmEy\nPON++fXdRIwG6ie41IkCr6AiuFpbzGBznqf+KyAveUYre0xnl0Uhf5aSIwkw\nuw5inZM/Mdr/ELNjf1tZH0yiohV7qCU6TdEmm5o9NLl/w7wQOeTLVWf+DC+h\nKxtCLSauEvb95IIt5TUpjgzyypmPNvkTAHfrNAo48qn5VvH4rUpQae1sLxln\n4dxC\r\n=K0Y5\r\n-----END PGP SIGNATURE-----\r\n","signatures":[{"keyid":"SHA256:jl3bwswu80PjjokCgh0o2w5c2U4LhQAE57gj9cz1kzA","sig":"MEUCIFbbzV9VW0pzbeo9puzDfeDIOzoum4f3+65nE+lBXu1QAiEA5zo2l/ndYY041N9uPfan9myJaqgmqnskrb1PxVOmZ6U="}]},"_npmUser":{"name":"danhawkins","email":"danny@quiqup.com"},"directories":{},"maintainers":[{"name":"danhawkins","email":"danny@quiqup.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages","tmp":"tmp/avro-kafkajs_0.8.2_1632936724894_0.3898643767560612"},"_hasShrinkwrap":false}},"time":{"created":"2021-09-29T17:27:45.169Z","0.8.0":"2021-09-29T17:27:45.414Z","modified":"2022-04-05T02:40:21.467Z","0.8.1":"2021-09-29T17:31:40.201Z","0.8.2":"2021-09-29T17:32:05.132Z"},"maintainers":[{"name":"danhawkins","email":"danny@quiqup.com"}],"description":"A wrapper around Kafkajs to transparently use Schema Registry for producing and consuming messages with avro schemas.","homepage":"https://github.com/ovotech/castle/tree/main/packages/avro-kafkajs#readme","repository":{"type":"git","url":"git@github.com:ovotech/castle.git"},"author":{"name":"Ivan Kerin","email":"ikerin@gmail.com"},"license":"Apache-2.0","readme":"# Avro Kafkajs\n\nA wrapper around [Kafka.js](https://github.com/tulios/kafkajs) to transparently use [Schema Registry](https://www.confluent.io/confluent-schema-registry/) for producing and consuming messages with [Avro schema](https://en.wikipedia.org/wiki/Apache_Avro).\n\n### Usage\n\n```shell\nyarn add @danhawkins/avro-kafkajs\n```\n\n> [examples/class.ts](examples/class.ts)\n\n```typescript\nimport { Kafka } from 'kafkajs';\nimport { SchemaRegistry, AvroKafka } from '@danhawkins/avro-kafkajs';\nimport { Schema } from 'avsc';\n\nconst mySchema: Schema = {\n  type: 'record',\n  name: 'MyMessage',\n  fields: [{ name: 'field1', type: 'string' }],\n};\n\n// Typescript types for the schema\ninterface MyMessage {\n  field1: string;\n}\n\nconst main = async () => {\n  const schemaRegistry = new SchemaRegistry({ uri: 'http://localhost:8081' });\n  const kafka = new Kafka({ brokers: ['localhost:29092'] });\n  const avroKafka = new AvroKafka(schemaRegistry, kafka);\n\n  // Consuming\n  const consumer = avroKafka.consumer({ groupId: 'my-group' });\n  await consumer.connect();\n  await consumer.subscribe({ topic: 'my-topic' });\n  await consumer.run<MyMessage>({\n    eachMessage: async ({ message }) => {\n      console.log(message.value);\n    },\n  });\n\n  // Producing\n  const producer = avroKafka.producer();\n  await producer.connect();\n  await producer.send<MyMessage>({\n    topic: 'my-topic',\n    schema: mySchema,\n    messages: [{ value: { field1: 'my-string' } }],\n  });\n};\n\nmain();\n```\n\nIt is a wrapper around [Kafka.js](https://github.com/tulios/kafkajs) with all of its functionality as is. With the one addition of requiring an `schema` field for the [Avro schema](https://en.wikipedia.org/wiki/Apache_Avro) when sending messages. Decoding of messages when consuming messages or batches happens automatically.\n\n### Encoded keys\n\nEncoding keys with avro is also supported:\n\n> [examples/encoded-key.ts](examples/encoded-key.ts)\n\n```typescript\nimport { Kafka } from 'kafkajs';\nimport { SchemaRegistry, AvroKafka } from '@danhawkins/avro-kafkajs';\nimport { Schema } from 'avsc';\n\nconst myValueSchema: Schema = {\n  type: 'record',\n  name: 'MyMessage',\n  fields: [{ name: 'field1', type: 'string' }],\n};\n\nconst myKeySchema: Schema = {\n  type: 'record',\n  name: 'MyKey',\n  fields: [{ name: 'id', type: 'int' }],\n};\n\n// Typescript types for the value schema\ninterface MyMessage {\n  field1: string;\n}\n\n// Typescript types for the key schema\ninterface MyKey {\n  id: number;\n}\n\nconst main = async () => {\n  const schemaRegistry = new SchemaRegistry({ uri: 'http://localhost:8081' });\n  const kafka = new Kafka({ brokers: ['localhost:29092'] });\n  const avroKafka = new AvroKafka(schemaRegistry, kafka);\n\n  // Consuming\n  const consumer = avroKafka.consumer({ groupId: 'my-group-key' });\n  await consumer.connect();\n  await consumer.subscribe({ topic: 'my-topic-with-key' });\n\n  // You need to specify that the key is encoded,\n  // otherwise it would just be returned as Buffer\n  // You can also pass the typescript type of the key\n  await consumer.run<MyMessage, MyKey>({\n    encodedKey: true,\n    eachMessage: async ({ message }) => {\n      console.log(message.key, message.value);\n    },\n  });\n\n  // Producing\n  const producer = avroKafka.producer();\n  await producer.connect();\n\n  // To produce messages, specify the keySchema and the key typescript type\n  await producer.send<MyMessage, MyKey>({\n    topic: 'my-topic-with-key',\n    schema: myValueSchema,\n    keySchema: myKeySchema,\n    messages: [{ value: { field1: 'my-string' }, key: { id: 111 } }],\n  });\n};\n\nmain();\n```\n\n### Using Schema Registry directly\n\nYou can also use schema registry directly to encode and decode messages.\n\n> [examples/schema-registry.ts](examples/schema-registry.ts)\n\n```typescript\nimport { Kafka } from 'kafkajs';\nimport { SchemaRegistry } from '@danhawkins/avro-kafkajs';\nimport { Schema } from 'avsc';\n\nconst mySchema: Schema = {\n  type: 'record',\n  name: 'MyMessage',\n  fields: [{ name: 'field1', type: 'string' }],\n};\n\nconst myKeySchema: Schema = {\n  type: 'record',\n  name: 'MyKey',\n  fields: [{ name: 'id', type: 'int' }],\n};\n\n// Typescript types for the schema\ninterface MyMessage {\n  field1: string;\n}\n\n// Typescript types for the key schema\ninterface MyKey {\n  id: number;\n}\n\nconst main = async () => {\n  const schemaRegistry = new SchemaRegistry({ uri: 'http://localhost:8081' });\n  const kafka = new Kafka({ brokers: ['localhost:29092'] });\n\n  // Consuming\n  const consumer = kafka.consumer({ groupId: 'my-group' });\n  await consumer.connect();\n  await consumer.subscribe({ topic: 'my-topic' });\n  await consumer.run({\n    eachMessage: async ({ message }) => {\n      const value = await schemaRegistry.decode<MyMessage>(message.value);\n      const key = await schemaRegistry.decode<MyKey>(message.key);\n      console.log(value, key);\n    },\n  });\n\n  // Producing\n  const producer = kafka.producer();\n  await producer.connect();\n\n  // Encode the value\n  const value = await schemaRegistry.encode<MyMessage>({\n    topic: 'my-topic',\n    schemaType: 'value',\n    schema: mySchema,\n    value: {\n      field1: 'my-string',\n    },\n  });\n\n  // Optionally encode the key\n  const key = await schemaRegistry.encode<MyKey>({\n    topic: 'my-topic',\n    schemaType: 'key',\n    schema: myKeySchema,\n    value: {\n      id: 10,\n    },\n  });\n  await producer.send({ topic: 'my-topic', messages: [{ value, key }] });\n};\n\nmain();\n```\n\n## Topic Aliases\n\nYou can define aliases to the topic names you want to listen to or produce messages for.\nThis can be used to encapsulate the real topic names, and use compile-time checked alias names throught your code.\n\n> [examples/topics-aliases.ts](examples/topics-aliases.ts)\n\n```typescript\nimport { Kafka } from 'kafkajs';\nimport { SchemaRegistry, AvroKafka, AvroProducer } from '@danhawkins/avro-kafkajs';\nimport { Schema } from 'avsc';\n\nconst mySchema: Schema = {\n  type: 'record',\n  name: 'MyMessage',\n  fields: [{ name: 'field1', type: 'string' }],\n};\n\n// Typescript types for the schema\ninterface MyMessage {\n  field1: string;\n}\n\nconst MY_TOPIC = 'myTopic';\n\n// Statically define a producer that would send the correct message to the correct topic\nconst sendMyMessage = (producer: AvroProducer, message: MyMessage) =>\n  producer.send<MyMessage>({\n    topic: MY_TOPIC,\n    schema: mySchema,\n    messages: [{ value: message, key: null }],\n  });\n\nconst main = async () => {\n  const aliases = { [MY_TOPIC]: 'my-topic-long-v1' };\n  const schemaRegistry = new SchemaRegistry({ uri: 'http://localhost:8081' });\n  const kafka = new Kafka({ brokers: ['localhost:29092'] });\n  const avroKafka = new AvroKafka(schemaRegistry, kafka, aliases);\n\n  // Consuming\n  const consumer = avroKafka.consumer({ groupId: 'my-group' });\n  await consumer.connect();\n  await consumer.subscribe({ topic: MY_TOPIC });\n  await consumer.run<MyMessage>({\n    eachMessage: async ({ message }) => {\n      console.log(message.value);\n    },\n  });\n\n  // Producing\n  const producer = avroKafka.producer();\n  await producer.connect();\n  await sendMyMessage(producer, { field1: 'my-string' });\n};\n\nmain();\n```\n\n## Schema evolution / Multiple schemas per topic\n\nWe can easily produce / consume different schemas for the same topic as the encoding / decoding logic will take it into account\n\n> [examples/schema-evolution.ts](examples/schema-evolution.ts)\n\n```typescript\nimport { Kafka } from 'kafkajs';\nimport { SchemaRegistry, AvroKafka } from '@danhawkins/avro-kafkajs';\nimport { Schema } from 'avsc';\n\nconst myOldSchema: Schema = {\n  type: 'record',\n  name: 'MyOldMessage',\n  fields: [{ name: 'field1', type: 'string' }],\n};\n\n// Backwards compatible schema change\nconst myNewSchema: Schema = {\n  type: 'record',\n  name: 'MyNewMessage',\n  fields: [\n    { name: 'field1', type: 'string' },\n    { name: 'field2', type: 'string', default: 'default-value' },\n  ],\n};\n\n// Typescript types for the schema\ninterface MyOldMessage {\n  field1: string;\n}\n\ninterface MyNewMessage {\n  field1: string;\n  field2?: string;\n}\n\nconst main = async () => {\n  const schemaRegistry = new SchemaRegistry({ uri: 'http://localhost:8081' });\n  const kafka = new Kafka({ brokers: ['localhost:29092'] });\n  const avroKafka = new AvroKafka(schemaRegistry, kafka);\n\n  // Consuming\n  const consumer = avroKafka.consumer({ groupId: 'my-group' });\n  await consumer.connect();\n  await consumer.subscribe({ topic: 'my-topic-evolution' });\n  await consumer.run<MyOldMessage | MyNewMessage>({\n    eachMessage: async ({ message }) => {\n      if ('field2' in message.value) {\n        // Typescript Type would match MyNewMessage\n        console.log('new message', message.value.field2);\n      } else {\n        // Typescript Type would match MyOldMessage\n        console.log('old message', message.value.field1);\n      }\n    },\n  });\n\n  // Producing\n  const producer = avroKafka.producer();\n  await producer.connect();\n  await producer.send<MyOldMessage>({\n    topic: 'my-topic-evolution',\n    schema: myOldSchema,\n    messages: [{ value: { field1: 'my-string' }, key: null }],\n  });\n  await producer.send<MyNewMessage>({\n    topic: 'my-topic-evolution',\n    schema: myNewSchema,\n    messages: [{ value: { field1: 'my-string', field2: 'new-string' }, key: null }],\n  });\n};\n\nmain();\n```\n\n## Custom schema registry subjects\n\nIf the subject for a topic in schema registry is already created, you can specify it directly to produce a message with the desired schema registry subject. It will take the latest version of the subject.\n\n> [examples/custom-subject.ts](examples/custom-subject.ts)\n\n```typescript\nimport { Kafka } from 'kafkajs';\nimport { SchemaRegistry, AvroKafka } from '@danhawkins/avro-kafkajs';\nimport { Schema } from 'avsc';\n\nconst mySchema: Schema = {\n  type: 'record',\n  name: 'MyMessage',\n  fields: [{ name: 'field1', type: 'string' }],\n};\n\n// Typescript types for the schema\ninterface MyMessage {\n  field1: string;\n}\n\nconst main = async () => {\n  const schemaRegistry = new SchemaRegistry({ uri: 'http://localhost:8081' });\n  const kafka = new Kafka({ brokers: ['localhost:29092'] });\n  const avroKafka = new AvroKafka(schemaRegistry, kafka);\n\n  // Consuming\n  const consumer = avroKafka.consumer({ groupId: 'my-group' });\n  await consumer.connect();\n  await consumer.subscribe({ topic: 'my-topic' });\n  await consumer.run<MyMessage>({\n    eachMessage: async ({ message }) => {\n      console.log(message.value);\n    },\n  });\n\n  // Producing\n  const producer = avroKafka.producer();\n  await producer.connect();\n  await producer.send<MyMessage>({\n    topic: 'my-topic',\n    schema: mySchema,\n    messages: [{ value: { field1: 'my-string' }, key: null }],\n  });\n\n  // Producing with custom subject\n  await producer.send<MyMessage>({\n    topic: 'my-topic',\n    subject: 'my-topic-value',\n    messages: [{ value: { field1: 'my-string-2' }, key: null }],\n  });\n};\n\nmain();\n```\n\n## Writing backfillers\n\nSometimes you'll want to write some code to backfill consumption using different data types, or test out consumption code. This package includes a transform stream to allow you to write a node stream -> batch payloads.\n\n> [examples/stream.ts](examples/stream.ts)\n\n```typescript\nimport { AvroTransformBatch } from '@danhawkins/avro-kafkajs';\nimport { Schema } from 'avsc';\nimport { ObjectReadableMock } from 'stream-mock';\n\nconst mySchema: Schema = {\n  type: 'record',\n  name: 'MyMessage',\n  fields: [{ name: 'field1', type: 'string' }],\n};\n\n// Typescript types for the schema\ninterface MyMessage {\n  field1: string;\n}\n\nconst data = new ObjectReadableMock(['one', 'two', 'three']);\n\nconst main = async () => {\n  const transform = new AvroTransformBatch<string, MyMessage, null>({\n    topic: 'test',\n    toKafkaMessage: (message) => ({\n      value: { field1: message },\n      key: null,\n      schema: mySchema,\n    }),\n  });\n\n  data.pipe(transform).on('data', (payload) => console.log(payload.batch.messages));\n};\n\nmain();\n```\n\n## Running the tests\n\nYou can run the tests with:\n\n```bash\nyarn test\n```\n\n### Coding style (linting, etc) tests\n\nStyle is maintained with prettier and eslint\n\n```\nyarn lint\n```\n\n## Deployment\n\nDeployment is preferment by lerna automatically on merge / push to main, but you'll need to bump the package version numbers yourself. Only updated packages with newer versions will be pushed to the npm registry.\n\n## Contributing\n\nHave a bug? File an issue with a simple example that reproduces this so we can take a look & confirm.\n\nWant to make a change? Submit a PR, explain why it's useful, and make sure you've updated the docs (this file) and the tests (see [test folder](test)).\n\n## License\n\nThis project is licensed under Apache 2 - see the [LICENSE](LICENSE) file for details\n","readmeFilename":"README.md"}