From 063908566cb59476420157b723d0068e092dd85c Mon Sep 17 00:00:00 2001 From: amiram Date: Tue, 18 Jul 2017 12:29:18 +0300 Subject: [PATCH 1/6] Add new KafkaClient more fixes to align with version 2 --- types/kafka-node/index.d.ts | 113 ++++++++++++++++++--------- types/kafka-node/kafka-node-tests.ts | 87 +++++++++++++++------ 2 files changed, 139 insertions(+), 61 deletions(-) diff --git a/types/kafka-node/index.d.ts b/types/kafka-node/index.d.ts index 06464d5602..ae1ced9f33 100644 --- a/types/kafka-node/index.d.ts +++ b/types/kafka-node/index.d.ts @@ -1,76 +1,73 @@ // Type definitions for kafka-node 1.3.3 // Project: https://github.com/SOHU-Co/kafka-node/ -// Definitions by: Daniel Imrie-Situnayake , Bill , Michael Haan +// Definitions by: Daniel Imrie-Situnayake , Bill , Michael Haan , Amiram Korach // Definitions: https://github.com/DefinitelyTyped/DefinitelyTyped - // # Classes export declare class Client { constructor(connectionString: string, clientId?: string, options?: ZKOptions, noBatchOptions?: AckBatchOptions, sslOptions?: any); close(callback?: Function): void; - topicExists(topics: Array, callback: Function): void; - refreshMetadata(topics: Array, cb?: (error: any, data: any) => any): void; - close(cb: (error: any) => any): void; - sendOffsetCommitV2Request(group: string, generationId: number, memberId: string, commits: Array, cb: (error: any, data: any) => any): void; + topicExists(topics: Array, callback: (err?: TopicsNotExistError | any) => any): void; + refreshMetadata(topics: Array, cb?: (error?: any) => any): void; + sendOffsetCommitV2Request(group: string, generationId: number, memberId: string, commits: Array, cb: Function): void; +} + +export declare class KafkaClient extends Client { + constructor(options?: KafkaClientOptions); + connect(): void; } export declare class Producer { - constructor(client: Client, options?: any, customPartitioner?: any); - on(eventName: string, cb: () => any): void; - on(eventName: string, cb: (error: any) => any): void; + constructor(client: Client, options?: ProducerOptions, customPartitioner?: any); + on(eventName: "ready", cb: () => any): void; + on(eventName: "error", cb: (error: any) => any): void; send(payloads: Array, cb: (error: any, data: any) => any): void; - createTopics(topics: Array, async: boolean, cb?: (error: any, data: any) => any): void; - close(cb: (error: any) => any): void; + createTopics(topics: Array, async: true, cb: (error?: any, data?: any) => any): void; + createTopics(topics: Array, async: false, cb: () => any): void; + createTopics(topics: Array, cb: (error?: any, data?: any) => any): void; + close(): void; } -export declare class HighLevelProducer { - constructor(client: Client, options?: any, customPartitioner?: any); - on(eventName: string, cb: () => any): void; - on(eventName: string, cb: (error: any) => any): void; - send(payloads: Array, cb: (error: any, data: any) => any): void; - createTopics(topics: Array, async: boolean, cb?: (error: any, data: any) => any): void; - close(cb: (error: any) => any): void; +export declare class HighLevelProducer extends Producer { } export declare class Consumer { - constructor(client: Client, fetchRequests: Array, options: ConsumerOptions); + constructor(client: Client, fetchRequests: Array, options: ConsumerOptions); client: Client; - on(eventName: string, cb: (message: string) => any): void; - on(eventName: string, cb: (error: any) => any): void; - addTopics(topics: Array, cb: (error: any, added: boolean) => any): void; - addTopics(topics: Array, cb: (error: any, added: boolean) => any, fromOffset: boolean): void; - removeTopics(topics: Array, cb: (error: any, removed: boolean) => any): void; + on(eventName: "message", cb: (message: Message) => any): void; + on(eventName: "error" | "offsetOutOfRange", cb: (error: any) => any): void; + addTopics(topics: Array | Array, cb: (error: any, added: Array | Array) => any): void; + addTopics(topics: Array | Array, cb: (error: any, added: Array | Array) => any, fromOffset: boolean): void; + removeTopics(topics: string | Array, cb: (error: any, removed: number) => any): void; commit(cb: (error: any, data: any) => any): void; + commit(force: boolean, cb: (error: any, data: any) => any): void; setOffset(topic: string, partition: number, offset: number): void; pause(): void; resume(): void; pauseTopics(topics: Array /* Array */): void; resumeTopics(topics: Array /* Array */): void; close(force: boolean, cb: () => any): void; + close(cb: () => any): void; } export declare class HighLevelConsumer { - constructor(client: Client, payloads: Array, options: ConsumerOptions); + constructor(client: Client, payloads: Array, options: HighLevelConsumerOptions); client: Client; - on(eventName: string, cb: (message: string) => any): void; - on(eventName: string, cb: (error: any) => any): void; - addTopics(topics: Array, cb: (error: any, added: boolean) => any): void; - addTopics(topics: Array, cb: (error: any, added: boolean) => any, fromOffset: boolean): void; - removeTopics(topics: Array, cb: (error: any, removed: boolean) => any): void; + on(eventName: "message", cb: (message: Message) => any): void; + on(eventName: "error" | "offsetOutOfRange", cb: (error: any) => any): void; + addTopics(topics: Array | Array, cb?: (error: any, added: Array | Array) => any): void; + removeTopics(topics: string | Array, cb: (error: any, removed: number) => any): void; commit(cb: (error: any, data: any) => any): void; + commit(force: boolean, cb: (error: any, data: any) => any): void; setOffset(topic: string, partition: number, offset: number): void; pause(): void; resume(): void; - pauseTopics(topics: Array /* Array */): void; - resumeTopics(topics: Array /* Array */): void; close(force: boolean, cb: () => any): void; + close(cb: () => any): void; } export declare class ConsumerGroup extends HighLevelConsumer { - constructor(options: ConsumerGroupOptions, topics: string[]); - on(eventName: string, cb: (message: string) => any): void; - on(eventName: string, cb: (error: any) => any): void; - close(force: boolean, cb: (error: any) => any): void; + constructor(options: ConsumerGroupOptions, topics: string[] | string); generationId: number; memberId: string; } @@ -91,6 +88,38 @@ export declare class KeyedMessage { } // # Interfaces + +export interface Message { + topic: string; + value: string; + offset: number; + partition: number; + highWaterOffset: number; + key: number; +} + +export interface ProducerOptions { + requireAcks?: number; + ackTimeoutMs?: number; + partitionerType?: number; +} + +export interface KafkaClientOptions { + kafkaHost?: string; + connectTimeout?: number; + requestTimeout?: number; + authConnect?: boolean; + connectRetryOptions?: ConnectRetryOptions; +} + +export interface ConnectRetryOptions { + retries?: number; + factor?: number; + minTimeout?: number; + maxTimeout?: number; + randomize?: boolean; +} + export interface AckBatchOptions { noAckBatchSize: number | null, noAckBatchAge: number | null @@ -112,7 +141,6 @@ export interface ProduceRequest { export interface ConsumerOptions { groupId?: string; - id?: string; autoCommit?: boolean; autoCommitIntervalMs?: number; fetchMaxWaitMs?: number; @@ -120,6 +148,11 @@ export interface ConsumerOptions { fetchMaxBytes?: number; fromOffset?: boolean; encoding?: string; + keyEncoding?: string; +} + +export interface HighLevelConsumerOptions extends ConsumerOptions { + id?: string; } export interface CustomPartitionAssignmentProtocol { @@ -134,7 +167,7 @@ export interface ConsumerGroupOptions { zk?: ZKOptions; batch?: AckBatchOptions; ssl?: boolean; - id: string; + id?: string; groupId: string; sessionTimeout?: number; protocol?: Array<"roundrobin" | "range" | CustomPartitionAssignmentProtocol>; @@ -180,3 +213,7 @@ export interface OffsetFetchRequest { partition?: number; offset?: number; } + +export declare class TopicsNotExistError extends Error { + topics: string | string[] +} diff --git a/types/kafka-node/kafka-node-tests.ts b/types/kafka-node/kafka-node-tests.ts index a348c8f771..8d99b98d63 100644 --- a/types/kafka-node/kafka-node-tests.ts +++ b/types/kafka-node/kafka-node-tests.ts @@ -12,9 +12,36 @@ var optionsClient = new kafka.Client('localhost:2181/', 'sendMessage', { }, { rejectUnauthorized: false }); +optionsClient.topicExists(['topic'], function (error: any){}); +optionsClient.refreshMetadata(['topic'], function (error: any){}); optionsClient.close(); +optionsClient.sendOffsetCommitV2Request('group', 0, 'memberId', [], function() {}); optionsClient.close(function(){}); +var basicKafkaClient = new kafka.KafkaClient(); + +var optionsKafkaClient = new kafka.KafkaClient({ + kafkaHost: 'localhost:2181', + connectTimeout: 1000, + requestTimeout: 1000, + authConnect: true, + connectRetryOptions: { + retries: 5, + factor: 0, + minTimeout: 1000, + maxTimeout: 1000, + randomize: true + } +}); + +optionsKafkaClient.connect(); + +var optionsProducer = new kafka.Producer(basicClient, { + requireAcks: 0, + ackTimeoutMs: 0, + partitionerType: 0 +}); + var producer = new kafka.Producer(basicClient); producer.on('error', function(error: Error){}); producer.on('ready', function(){ @@ -43,10 +70,10 @@ producer.on('ready', function(){ producer.send(messages, function(err: Error){}); producer.send(messages, function(err: Error, data: Object){}); - producer.createTopics(['t'], true, function (err: Error, data: Object) {}); - producer.createTopics(['t'], false, function (err, data) {}); - // producer.createTopics(['t'], function (err: Error, data: Object) {}); // Omitting middle argument is not possible in TS - + producer.createTopics(['t'], true, function (err: Error, data: any) {}); + producer.createTopics(['t'], function (err: Error, data: any) {}); + producer.createTopics(['t'], false, function () {}); + producer.close(); }); var highLevelProducer = new kafka.HighLevelProducer(basicClient); @@ -73,13 +100,13 @@ highLevelProducer.on('ready', function(){ messages: [new kafka.KeyedMessage('key', 'message')] }]; - producer.send(messages, function(err: Error){}); - producer.send(messages, function(err: Error, data: Object){}); - - producer.createTopics(['t'], true, function (err: Error, data: Object) {}); - producer.createTopics(['t'], false, function (err, data) {}); - // producer.createTopics(['t'], function (err: Error, data: Object) {}); // Omitting middle argument is not possible in TS + highLevelProducer.send(messages, function(err: Error){}); + highLevelProducer.send(messages, function(err: Error, data: Object){}); + producer.createTopics(['t'], true, function (err: Error, data: any) {}); + producer.createTopics(['t'], function (err: Error, data: any) {}); + producer.createTopics(['t'], false, function () {}); + producer.close(); }); var fetchRequests = [{ topic: 'awesome' }]; @@ -88,14 +115,24 @@ var consumer = new kafka.Consumer(basicClient, fetchRequests, { autoCommit: true }); consumer.on('error', function(error: Error){}); -consumer.on('message', function(message){}); +consumer.on('offsetOutOfRange', function(error: Error){}); +consumer.on('message', function(message: kafka.Message){ + var topic = message.topic; + var value = message.value; + var offset = message.offset; + var partition = message.partition; + var highWaterOffset = message.highWaterOffset; + var key = message.key; +}); consumer.addTopics(['t1', 't2'], function (err, added) {}); consumer.addTopics([{ topic: 't1', offset: 10 }], function (err, added) {}, true); -consumer.removeTopics(['t1', 't2'], function (err, removed) {}); +consumer.removeTopics(['t1', 't2'], function (err, removed: number) {}); +consumer.removeTopics('t2', function (err, removed: number) {}); consumer.commit(function (err, data) {}); +consumer.commit(true, function (err, data) {}); consumer.setOffset('topic', 0, 0); @@ -111,6 +148,7 @@ consumer.resumeTopics([ ]); consumer.close(true, function () {}); +consumer.close(function () {}); var fetchRequests = [{ topic: 'awesome' }]; var hlConsumer = new kafka.HighLevelConsumer(basicClient, fetchRequests, { @@ -119,28 +157,31 @@ var hlConsumer = new kafka.HighLevelConsumer(basicClient, fetchRequests, { }); hlConsumer.on('error', function(error: Error){}); -hlConsumer.on('message', function(message){}); +hlConsumer.on('offsetOutOfRange', function(error: Error){}); +hlConsumer.on('message', function(message: kafka.Message){ + var topic = message.topic; + var value = message.value; + var offset = message.offset; + var partition = message.partition; + var highWaterOffset = message.highWaterOffset; + var key = message.key; +}); hlConsumer.addTopics(['t1', 't2'], function (err, added) {}); -hlConsumer.addTopics([{ topic: 't1', offset: 10 }], function (err, added) {}, true); +hlConsumer.addTopics([{ topic: 't1', offset: 10 }], function (err, added) {}); -hlConsumer.removeTopics(['t1', 't2'], function (err, removed) {}); +hlConsumer.removeTopics(['t1', 't2'], function (err, removed: number) {}); +hlConsumer.removeTopics('t2', function (err, removed: number) {}); hlConsumer.commit(function (err, data) {}); +hlConsumer.commit(true, function (err, data) {}); hlConsumer.setOffset('topic', 0, 0); hlConsumer.pause(); hlConsumer.resume(); -hlConsumer.pauseTopics([ - 'topic1', - { topic: 'topic2', partition: 0 } -]); -hlConsumer.resumeTopics([ - 'topic1', - { topic: 'topic2', partition: 0 } -]); hlConsumer.close(true, function () {}); +hlConsumer.close(function () {}); var ackBatchOptions = {'noAckBatchSize': 1024, 'noAckBatchAge': 10}; var cgOptions: kafka.ConsumerGroupOptions = { From cc5bb4de9396a658ff12fdaf221cfec3932b3888 Mon Sep 17 00:00:00 2001 From: amiram Date: Tue, 18 Jul 2017 12:31:25 +0300 Subject: [PATCH 2/6] message properties optional --- types/kafka-node/index.d.ts | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/types/kafka-node/index.d.ts b/types/kafka-node/index.d.ts index ae1ced9f33..5fc3ef958e 100644 --- a/types/kafka-node/index.d.ts +++ b/types/kafka-node/index.d.ts @@ -92,10 +92,10 @@ export declare class KeyedMessage { export interface Message { topic: string; value: string; - offset: number; - partition: number; - highWaterOffset: number; - key: number; + offset?: number; + partition?: number; + highWaterOffset?: number; + key?: number; } export interface ProducerOptions { From 45cc08937da9f3d8f8524c050f97790a5fc5ad68 Mon Sep 17 00:00:00 2001 From: amiram Date: Tue, 18 Jul 2017 12:33:54 +0300 Subject: [PATCH 3/6] add tslint --- types/kafka-node/tsconfig.json | 4 ++-- types/kafka-node/tslint.json | 1 + 2 files changed, 3 insertions(+), 2 deletions(-) create mode 100644 types/kafka-node/tslint.json diff --git a/types/kafka-node/tsconfig.json b/types/kafka-node/tsconfig.json index f049d379df..616ae7507a 100644 --- a/types/kafka-node/tsconfig.json +++ b/types/kafka-node/tsconfig.json @@ -6,7 +6,7 @@ ], "noImplicitAny": true, "noImplicitThis": true, - "strictNullChecks": false, + "strictNullChecks": true, "baseUrl": "../", "typeRoots": [ "../" @@ -19,4 +19,4 @@ "index.d.ts", "kafka-node-tests.ts" ] -} \ No newline at end of file +} diff --git a/types/kafka-node/tslint.json b/types/kafka-node/tslint.json new file mode 100644 index 0000000000..0967ef424b --- /dev/null +++ b/types/kafka-node/tslint.json @@ -0,0 +1 @@ +{} From 34f5e746e2a25d0769ce0c7bffb5768b06957b30 Mon Sep 17 00:00:00 2001 From: amiram Date: Tue, 18 Jul 2017 13:05:42 +0300 Subject: [PATCH 4/6] fix lint --- types/kafka-node/index.d.ts | 76 +++--- types/kafka-node/kafka-node-tests.ts | 335 +++++++++++++++------------ types/kafka-node/tslint.json | 4 +- 3 files changed, 229 insertions(+), 186 deletions(-) diff --git a/types/kafka-node/index.d.ts b/types/kafka-node/index.d.ts index 5fc3ef958e..5c93822a03 100644 --- a/types/kafka-node/index.d.ts +++ b/types/kafka-node/index.d.ts @@ -1,62 +1,61 @@ -// Type definitions for kafka-node 1.3.3 +// Type definitions for kafka-node 2.0 // Project: https://github.com/SOHU-Co/kafka-node/ // Definitions by: Daniel Imrie-Situnayake , Bill , Michael Haan , Amiram Korach // Definitions: https://github.com/DefinitelyTyped/DefinitelyTyped // # Classes -export declare class Client { +export class Client { constructor(connectionString: string, clientId?: string, options?: ZKOptions, noBatchOptions?: AckBatchOptions, sslOptions?: any); - close(callback?: Function): void; - topicExists(topics: Array, callback: (err?: TopicsNotExistError | any) => any): void; - refreshMetadata(topics: Array, cb?: (error?: any) => any): void; - sendOffsetCommitV2Request(group: string, generationId: number, memberId: string, commits: Array, cb: Function): void; + close(callback?: () => void): void; + topicExists(topics: string[], callback: (err?: TopicsNotExistError | any) => any): void; + refreshMetadata(topics: string[], cb?: (error?: any) => any): void; + sendOffsetCommitV2Request(group: string, generationId: number, memberId: string, commits: OffsetCommitRequest[], cb: () => void): void; } -export declare class KafkaClient extends Client { +export class KafkaClient extends Client { constructor(options?: KafkaClientOptions); connect(): void; } -export declare class Producer { +export class Producer { constructor(client: Client, options?: ProducerOptions, customPartitioner?: any); on(eventName: "ready", cb: () => any): void; on(eventName: "error", cb: (error: any) => any): void; - send(payloads: Array, cb: (error: any, data: any) => any): void; - createTopics(topics: Array, async: true, cb: (error?: any, data?: any) => any): void; - createTopics(topics: Array, async: false, cb: () => any): void; - createTopics(topics: Array, cb: (error?: any, data?: any) => any): void; + send(payloads: ProduceRequest[], cb: (error: any, data: any) => any): void; + createTopics(topics: string[], async: true, cb: (error?: any, data?: any) => any): void; + createTopics(topics: string[], async: false, cb: () => any): void; + createTopics(topics: string[], cb: (error?: any, data?: any) => any): void; close(): void; } -export declare class HighLevelProducer extends Producer { +export class HighLevelProducer extends Producer { } -export declare class Consumer { +export class Consumer { constructor(client: Client, fetchRequests: Array, options: ConsumerOptions); client: Client; on(eventName: "message", cb: (message: Message) => any): void; on(eventName: "error" | "offsetOutOfRange", cb: (error: any) => any): void; - addTopics(topics: Array | Array, cb: (error: any, added: Array | Array) => any): void; - addTopics(topics: Array | Array, cb: (error: any, added: Array | Array) => any, fromOffset: boolean): void; - removeTopics(topics: string | Array, cb: (error: any, removed: number) => any): void; + addTopics(topics: string[] | Topic[], cb: (error: any, added: string[] | Topic[]) => any, fromOffset?: boolean): void; + removeTopics(topics: string | string[], cb: (error: any, removed: number) => any): void; commit(cb: (error: any, data: any) => any): void; commit(force: boolean, cb: (error: any, data: any) => any): void; setOffset(topic: string, partition: number, offset: number): void; pause(): void; resume(): void; - pauseTopics(topics: Array /* Array */): void; - resumeTopics(topics: Array /* Array */): void; + pauseTopics(topics: any[] /* Array */): void; + resumeTopics(topics: any[] /* Array */): void; close(force: boolean, cb: () => any): void; close(cb: () => any): void; } -export declare class HighLevelConsumer { - constructor(client: Client, payloads: Array, options: HighLevelConsumerOptions); +export class HighLevelConsumer { + constructor(client: Client, payloads: Topic[], options: HighLevelConsumerOptions); client: Client; on(eventName: "message", cb: (message: Message) => any): void; on(eventName: "error" | "offsetOutOfRange", cb: (error: any) => any): void; - addTopics(topics: Array | Array, cb?: (error: any, added: Array | Array) => any): void; - removeTopics(topics: string | Array, cb: (error: any, removed: number) => any): void; + addTopics(topics: string[] | Topic[], cb?: (error: any, added: string[] | Topic[]) => any): void; + removeTopics(topics: string | string[], cb: (error: any, removed: number) => any): void; commit(cb: (error: any, data: any) => any): void; commit(force: boolean, cb: (error: any, data: any) => any): void; setOffset(topic: string, partition: number, offset: number): void; @@ -66,24 +65,23 @@ export declare class HighLevelConsumer { close(cb: () => any): void; } -export declare class ConsumerGroup extends HighLevelConsumer { +export class ConsumerGroup extends HighLevelConsumer { constructor(options: ConsumerGroupOptions, topics: string[] | string); generationId: number; memberId: string; } -export declare class Offset { +export class Offset { constructor(client: Client); - on(eventName: string, cb: () => any): void; - fetch(payloads: Array, cb: (error: any, data: any) => any): void; - commit(groupId: string, payloads: Array, cb: (error: any, data: any) => any): void; - fetchCommits(groupId: string, payloads: Array, cb: (error: any, data: any) => any): void; - fetchLatestOffsets(topics: Array, cb: (error: any, data: any) => any): void; - fetchEarliestOffsets(topics: Array, cb: (error: any, data: any) => any): void; - on(eventName: string, cb: (error: any) => any): void; + on(eventName: string, cb: (error?: any) => any): void; + fetch(payloads: OffsetRequest[], cb: (error: any, data: any) => any): void; + commit(groupId: string, payloads: OffsetCommitRequest[], cb: (error: any, data: any) => any): void; + fetchCommits(groupId: string, payloads: OffsetFetchRequest[], cb: (error: any, data: any) => any): void; + fetchLatestOffsets(topics: string[], cb: (error: any, data: any) => any): void; + fetchEarliestOffsets(topics: string[], cb: (error: any, data: any) => any): void; } -export declare class KeyedMessage { +export class KeyedMessage { constructor(key: string, message: string); } @@ -121,8 +119,8 @@ export interface ConnectRetryOptions { } export interface AckBatchOptions { - noAckBatchSize: number | null, - noAckBatchAge: number | null + noAckBatchSize: number | null; + noAckBatchAge: number | null; } export interface ZKOptions { @@ -133,7 +131,7 @@ export interface ZKOptions { export interface ProduceRequest { topic: string; - messages: any; // Array | Array | string | KeyedMessage + messages: any; // string[] | Array | string | KeyedMessage key?: any; partition?: number; attributes?: number; @@ -159,7 +157,7 @@ export interface CustomPartitionAssignmentProtocol { name: string; version: number; userData: {}; - assign: (topicPattern: any, groupMembers: any, callback: (error: any, result: any) => void) => void; + assign(topicPattern: any, groupMembers: any, callback: (error: any, result: any) => void): void; } export interface ConsumerGroupOptions { @@ -214,6 +212,6 @@ export interface OffsetFetchRequest { offset?: number; } -export declare class TopicsNotExistError extends Error { - topics: string | string[] +export class TopicsNotExistError extends Error { + topics: string | string[]; } diff --git a/types/kafka-node/kafka-node-tests.ts b/types/kafka-node/kafka-node-tests.ts index 8d99b98d63..8a3553a9c5 100644 --- a/types/kafka-node/kafka-node-tests.ts +++ b/types/kafka-node/kafka-node-tests.ts @@ -1,26 +1,30 @@ import kafka = require('kafka-node'); -var basicClient = new kafka.Client('localhost:2181/', 'sendMessage'); +const basicClient = new kafka.Client('localhost:2181/', 'sendMessage'); -var optionsClient = new kafka.Client('localhost:2181/', 'sendMessage', { - sessionTimeout: 30000, - spinDelay: 1000, - retries: 0 +const optionsClient = new kafka.Client('localhost:2181/', 'sendMessage', { + sessionTimeout: 30000, + spinDelay: 1000, + retries: 0 }, { noAckBatchSize: 1000, noAckBatchAge: 1000 * 10 }, { rejectUnauthorized: false }); -optionsClient.topicExists(['topic'], function (error: any){}); -optionsClient.refreshMetadata(['topic'], function (error: any){}); +optionsClient.topicExists(['topic'], (error: any) => { +}); +optionsClient.refreshMetadata(['topic'], (error: any) => { +}); optionsClient.close(); -optionsClient.sendOffsetCommitV2Request('group', 0, 'memberId', [], function() {}); -optionsClient.close(function(){}); +optionsClient.sendOffsetCommitV2Request('group', 0, 'memberId', [], () => { +}); +optionsClient.close(() => { +}); -var basicKafkaClient = new kafka.KafkaClient(); +const basicKafkaClient = new kafka.KafkaClient(); -var optionsKafkaClient = new kafka.KafkaClient({ +const optionsKafkaClient = new kafka.KafkaClient({ kafkaHost: 'localhost:2181', connectTimeout: 1000, requestTimeout: 1000, @@ -36,187 +40,226 @@ var optionsKafkaClient = new kafka.KafkaClient({ optionsKafkaClient.connect(); -var optionsProducer = new kafka.Producer(basicClient, { +const optionsProducer = new kafka.Producer(basicClient, { requireAcks: 0, ackTimeoutMs: 0, partitionerType: 0 }); -var producer = new kafka.Producer(basicClient); -producer.on('error', function(error: Error){}); -producer.on('ready', function(){ - - var messages = [{ - topic: 'topicName', - messages: ['message body'], - partition: 0, - attributes: 2 - }, { - topic: 'topicName', - messages: ['message body'], - partition: 0 - }, { - topic: 'topicName', - messages: ['message body'], - attributes: 0 - }, { - topic: 'topicName', - messages: ['message body'] - }, { - topic: 'topicName', - messages: [new kafka.KeyedMessage('key', 'message')] - }]; - - producer.send(messages, function(err: Error){}); - producer.send(messages, function(err: Error, data: Object){}); - - producer.createTopics(['t'], true, function (err: Error, data: any) {}); - producer.createTopics(['t'], function (err: Error, data: any) {}); - producer.createTopics(['t'], false, function () {}); - producer.close(); +const producer = new kafka.Producer(basicClient); +producer.on('error', (error: Error) => { }); +producer.on('ready', () => { + const messages = [{ + topic: 'topicName', + messages: ['message body'], + partition: 0, + attributes: 2 + }, { + topic: 'topicName', + messages: ['message body'], + partition: 0 + }, { + topic: 'topicName', + messages: ['message body'], + attributes: 0 + }, { + topic: 'topicName', + messages: ['message body'] + }, { + topic: 'topicName', + messages: [new kafka.KeyedMessage('key', 'message')] + }]; -var highLevelProducer = new kafka.HighLevelProducer(basicClient); -highLevelProducer.on('error', function(error: Error){}); -highLevelProducer.on('ready', function(){ + producer.send(messages, (err: Error) => { + }); + producer.send(messages, (err: Error, data: any) => { + }); - var messages = [{ - topic: 'topicName', - messages: ['message body'], - attributes: 2 - }, { - topic: 'topicName', - messages: ['message body'], - partition: 0 - }, { - topic: 'topicName', - messages: ['message body'], - attributes: 0 - }, { - topic: 'topicName', - messages: ['message body'] - }, { - topic: 'topicName', - messages: [new kafka.KeyedMessage('key', 'message')] - }]; - - highLevelProducer.send(messages, function(err: Error){}); - highLevelProducer.send(messages, function(err: Error, data: Object){}); - - producer.createTopics(['t'], true, function (err: Error, data: any) {}); - producer.createTopics(['t'], function (err: Error, data: any) {}); - producer.createTopics(['t'], false, function () {}); + producer.createTopics(['t'], true, (err: Error, data: any) => { + }); + producer.createTopics(['t'], (err: Error, data: any) => { + }); + producer.createTopics(['t'], false, () => { + }); producer.close(); }); -var fetchRequests = [{ topic: 'awesome' }]; -var consumer = new kafka.Consumer(basicClient, fetchRequests, { - groupId: 'abcde', - autoCommit: true +const highLevelProducer = new kafka.HighLevelProducer(basicClient); +highLevelProducer.on('error', (error: Error) => { }); -consumer.on('error', function(error: Error){}); -consumer.on('offsetOutOfRange', function(error: Error){}); -consumer.on('message', function(message: kafka.Message){ - var topic = message.topic; - var value = message.value; - var offset = message.offset; - var partition = message.partition; - var highWaterOffset = message.highWaterOffset; - var key = message.key; +highLevelProducer.on('ready', () => { + const messages = [{ + topic: 'topicName', + messages: ['message body'], + attributes: 2 + }, { + topic: 'topicName', + messages: ['message body'], + partition: 0 + }, { + topic: 'topicName', + messages: ['message body'], + attributes: 0 + }, { + topic: 'topicName', + messages: ['message body'] + }, { + topic: 'topicName', + messages: [new kafka.KeyedMessage('key', 'message')] + }]; + + highLevelProducer.send(messages, (err: Error) => { + }); + highLevelProducer.send(messages, (err: Error, data: any) => { + }); + + producer.createTopics(['t'], true, (err: Error, data: any) => { + }); + producer.createTopics(['t'], (err: Error, data: any) => { + }); + producer.createTopics(['t'], false, () => { + }); + producer.close(); }); -consumer.addTopics(['t1', 't2'], function (err, added) {}); -consumer.addTopics([{ topic: 't1', offset: 10 }], function (err, added) {}, true); +const fetchRequests = [{topic: 'awesome'}]; +const consumer = new kafka.Consumer(basicClient, fetchRequests, { + groupId: 'abcde', + autoCommit: true +}); +consumer.on('error', (error: Error) => { +}); +consumer.on('offsetOutOfRange', (error: Error) => { +}); +consumer.on('message', (message: kafka.Message) => { + const topic = message.topic; + const value = message.value; + const offset = message.offset; + const partition = message.partition; + const highWaterOffset = message.highWaterOffset; + const key = message.key; +}); -consumer.removeTopics(['t1', 't2'], function (err, removed: number) {}); -consumer.removeTopics('t2', function (err, removed: number) {}); +consumer.addTopics(['t1', 't2'], (err: any, added: any) => { +}); +consumer.addTopics([{topic: 't1', offset: 10}], (err: any, added: any) => { +}, true); -consumer.commit(function (err, data) {}); -consumer.commit(true, function (err, data) {}); +consumer.removeTopics(['t1', 't2'], (err: any, removed: number) => { +}); +consumer.removeTopics('t2', (err: any, removed: number) => { +}); + +consumer.commit((err: any, data: any) => { +}); +consumer.commit(true, (err: any, data: any) => { +}); consumer.setOffset('topic', 0, 0); consumer.pause(); consumer.resume(); consumer.pauseTopics([ - 'topic1', - { topic: 'topic2', partition: 0 } + 'topic1', + {topic: 'topic2', partition: 0} ]); consumer.resumeTopics([ - 'topic1', - { topic: 'topic2', partition: 0 } + 'topic1', + {topic: 'topic2', partition: 0} ]); -consumer.close(true, function () {}); -consumer.close(function () {}); - -var fetchRequests = [{ topic: 'awesome' }]; -var hlConsumer = new kafka.HighLevelConsumer(basicClient, fetchRequests, { - groupId: 'abcde', - autoCommit: true +consumer.close(true, () => { +}); +consumer.close(() => { }); -hlConsumer.on('error', function(error: Error){}); -hlConsumer.on('offsetOutOfRange', function(error: Error){}); -hlConsumer.on('message', function(message: kafka.Message){ - var topic = message.topic; - var value = message.value; - var offset = message.offset; - var partition = message.partition; - var highWaterOffset = message.highWaterOffset; - var key = message.key; +const fetchRequests1 = [{topic: 'awesome'}]; +const hlConsumer = new kafka.HighLevelConsumer(basicClient, fetchRequests1, { + groupId: 'abcde', + autoCommit: true }); -hlConsumer.addTopics(['t1', 't2'], function (err, added) {}); -hlConsumer.addTopics([{ topic: 't1', offset: 10 }], function (err, added) {}); -hlConsumer.removeTopics(['t1', 't2'], function (err, removed: number) {}); -hlConsumer.removeTopics('t2', function (err, removed: number) {}); +hlConsumer.on('error', (error: Error) => { +}); +hlConsumer.on('offsetOutOfRange', (error: Error) => { +}); +hlConsumer.on('message', (message: kafka.Message) => { + const topic = message.topic; + const value = message.value; + const offset = message.offset; + const partition = message.partition; + const highWaterOffset = message.highWaterOffset; + const key = message.key; +}); +hlConsumer.addTopics(['t1', 't2'], (err: any, added: any) => { +}); +hlConsumer.addTopics([{topic: 't1', offset: 10}], (err: any, added: any) => { +}); -hlConsumer.commit(function (err, data) {}); -hlConsumer.commit(true, function (err, data) {}); +hlConsumer.removeTopics(['t1', 't2'], (err: any, removed: number) => { +}); +hlConsumer.removeTopics('t2', (err: any, removed: number) => { +}); + +hlConsumer.commit((err: any, data: any) => { +}); +hlConsumer.commit(true, (err: any, data: any) => { +}); hlConsumer.setOffset('topic', 0, 0); hlConsumer.pause(); hlConsumer.resume(); -hlConsumer.close(true, function () {}); -hlConsumer.close(function () {}); +hlConsumer.close(true, () => { +}); +hlConsumer.close(() => { +}); -var ackBatchOptions = {'noAckBatchSize': 1024, 'noAckBatchAge': 10}; -var cgOptions: kafka.ConsumerGroupOptions = { - host: 'localhost:2181/', - batch: ackBatchOptions, - groupId: 'groupID', - id: 'consumerID', - sessionTimeout: 15000, - protocol: ["roundrobin"], - fromOffset: "latest", - migrateHLC: false, - migrateRolling: true +const ackBatchOptions = {noAckBatchSize: 1024, noAckBatchAge: 10}; +const cgOptions: kafka.ConsumerGroupOptions = { + host: 'localhost:2181/', + batch: ackBatchOptions, + groupId: 'groupID', + id: 'consumerID', + sessionTimeout: 15000, + protocol: ["roundrobin"], + fromOffset: "latest", + migrateHLC: false, + migrateRolling: true }; -var consumerGroup = new kafka.ConsumerGroup( cgOptions, ['topic1']); -consumerGroup.on('error', (err) => {}); -consumerGroup.on('message', (msg) => {}); -consumerGroup.close(true, () => {}); +const consumerGroup = new kafka.ConsumerGroup(cgOptions, ['topic1']); +consumerGroup.on('error', (err) => { +}); +consumerGroup.on('message', (msg) => { +}); +consumerGroup.close(true, () => { +}); -var offset = new kafka.Offset(basicClient); +const offset = new kafka.Offset(basicClient); -offset.on('ready', function(){}); +offset.on('ready', () => { +}); offset.fetch([ - { topic: 't', partition: 0, time: Date.now(), maxNum: 1 }, - { topic: 't' } -], function (err, data) { }); + {topic: 't', partition: 0, time: Date.now(), maxNum: 1}, + {topic: 't'} +], (err: any, data: any) => { +}); offset.commit('groupId', [ - { topic: 't', partition: 0, offset: 10 } -], function (err, data) { }); + {topic: 't', partition: 0, offset: 10} +], (err, data) => { +}); offset.fetchCommits('groupId', [ - { topic: 't', partition: 0 } -], function (err, data) {}); + {topic: 't', partition: 0} +], (err, data) => { +}); -offset.fetchLatestOffsets(['t'], (err, offsets) => {}) -offset.fetchEarliestOffsets(['t'], (err, offsets) => {}) +offset.fetchLatestOffsets(['t'], (err, offsets) => { +}); +offset.fetchEarliestOffsets(['t'], (err, offsets) => { +}); diff --git a/types/kafka-node/tslint.json b/types/kafka-node/tslint.json index 0967ef424b..d88586e5bd 100644 --- a/types/kafka-node/tslint.json +++ b/types/kafka-node/tslint.json @@ -1 +1,3 @@ -{} +{ + "extends": "dtslint/dt.json" +} From aa245d6c8c1cebdb85ce5018892b5da5ca9316cc Mon Sep 17 00:00:00 2001 From: amiram Date: Tue, 18 Jul 2017 22:20:13 +0300 Subject: [PATCH 5/6] add missing kafkaClient options --- types/kafka-node/index.d.ts | 2 ++ 1 file changed, 2 insertions(+) diff --git a/types/kafka-node/index.d.ts b/types/kafka-node/index.d.ts index 5c93822a03..d025f849ae 100644 --- a/types/kafka-node/index.d.ts +++ b/types/kafka-node/index.d.ts @@ -108,6 +108,8 @@ export interface KafkaClientOptions { requestTimeout?: number; authConnect?: boolean; connectRetryOptions?: ConnectRetryOptions; + sslOptions?: any; + clientId?: string; } export interface ConnectRetryOptions { From 59932cbdd60a267ab18d479506c6108f48d58100 Mon Sep 17 00:00:00 2001 From: amiram Date: Tue, 18 Jul 2017 22:21:14 +0300 Subject: [PATCH 6/6] add test for missing kafkaClient options --- types/kafka-node/kafka-node-tests.ts | 2 ++ 1 file changed, 2 insertions(+) diff --git a/types/kafka-node/kafka-node-tests.ts b/types/kafka-node/kafka-node-tests.ts index 8a3553a9c5..b1f10eb291 100644 --- a/types/kafka-node/kafka-node-tests.ts +++ b/types/kafka-node/kafka-node-tests.ts @@ -29,6 +29,8 @@ const optionsKafkaClient = new kafka.KafkaClient({ connectTimeout: 1000, requestTimeout: 1000, authConnect: true, + sslOptions: {}, + clientId: "client id", connectRetryOptions: { retries: 5, factor: 0,