From 4c8ac68e402069b7b0ce06f6455618d8d888caa9 Mon Sep 17 00:00:00 2001 From: Alexander T Date: Mon, 4 Nov 2019 23:42:29 +0200 Subject: [PATCH] kafka-node: Provides its own types (#39580) --- notNeededPackages.json | 6 + types/kafka-node/index.d.ts | 271 ------------------------ types/kafka-node/kafka-node-tests.ts | 298 --------------------------- types/kafka-node/tsconfig.json | 23 --- types/kafka-node/tslint.json | 8 - 5 files changed, 6 insertions(+), 600 deletions(-) delete mode 100644 types/kafka-node/index.d.ts delete mode 100644 types/kafka-node/kafka-node-tests.ts delete mode 100644 types/kafka-node/tsconfig.json delete mode 100644 types/kafka-node/tslint.json diff --git a/notNeededPackages.json b/notNeededPackages.json index 0f98be0618..0105784226 100644 --- a/notNeededPackages.json +++ b/notNeededPackages.json @@ -2184,6 +2184,12 @@ "sourceRepoURL": "https://github.com/sindresorhus/junk", "asOfVersion": "3.0.0" }, + { + "libraryName": "kafka-node", + "typingsPackageName": "kafka-node", + "sourceRepoURL": "https://github.com/SOHU-Co/kafka-node/", + "asOfVersion": "3.0.0" + }, { "libraryName": "karma-viewport", "typingsPackageName": "karma-viewport", diff --git a/types/kafka-node/index.d.ts b/types/kafka-node/index.d.ts deleted file mode 100644 index 441f9e88c2..0000000000 --- a/types/kafka-node/index.d.ts +++ /dev/null @@ -1,271 +0,0 @@ -// 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 -// Insanehong -// Roger -// Definitions: https://github.com/DefinitelyTyped/DefinitelyTyped - -/// - -// # Classes -export class Client { - constructor(connectionString: string, clientId?: string, options?: ZKOptions, noBatchOptions?: AckBatchOptions, sslOptions?: any); - close(cb?: () => void): void; - loadMetadataForTopics(topics: string[], cb: (error: TopicsNotExistError | any, data: any) => any): void; - topicExists(topics: string[], cb: (error?: TopicsNotExistError | any) => any): void; - refreshMetadata(topics: string[], cb?: (error?: any) => any): void; - sendOffsetCommitV2Request(group: string, generationId: number, memberId: string, commits: OffsetCommitRequest[], cb: (error: any, data: any) => any): void; - // Note: socket_error is currently KafkaClient only, and zkReconnect is currently Client only. - on(eventName: "brokersChanged" | "close" | "connect" | "ready" | "reconnect" | "zkReconnect", cb: () => any): this; - on(eventName: "error" | "socket_error", cb: (error: any) => any): this; -} - -export class KafkaClient extends Client { - constructor(options?: KafkaClientOptions); - connect(): void; - getListGroups(cb: (error: any, data: any) => any): void; - describeGroups(consumerGroups: any, cb: (error: any, data: any) => any): void; -} - -export class Producer { - constructor(client: Client, options?: ProducerOptions, customPartitioner?: CustomPartitioner); - on(eventName: "ready", cb: () => any): void; - on(eventName: "error", cb: (error: any) => any): void; - send(payloads: ProduceRequest[], cb: (error: any, data: any) => any): void; - createTopics(topics: string[], async: boolean, cb: (error: any, data: any) => any): void; - createTopics(topics: string[], cb: (error: any, data: any) => any): void; - close(cb?: () => any): void; -} - -export class HighLevelProducer extends Producer { -} - -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: 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: any[] /* Array */): void; - resumeTopics(topics: any[] /* Array */): void; - close(force: boolean, cb: () => any): void; - close(cb: () => any): void; -} - -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; - on(eventName: "rebalancing" | "rebalanced" | "connect", cb: () => 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; - sendOffsetCommitRequest(commits: OffsetCommitRequest[], cb: (error: any, data: any) => any): void; - setOffset(topic: string, partition: number, offset: number): void; - pause(): void; - resume(): void; - close(force: boolean, cb: (error: any) => any): void; - close(cb: () => any): void; -} - -export class ConsumerGroup extends HighLevelConsumer { - constructor(options: ConsumerGroupOptions, topics: string[] | string); - generationId: number; - memberId: string; - client: KafkaClient & Client; -} - -export class Offset { - constructor(client: Client); - on(eventName: "ready" | "connect", cb: () => any): void; - on(eventName: "error", 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 class KeyedMessage { - constructor(key: string, value: string | Buffer); -} - -export class Admin { - constructor(kafkaClient: KafkaClient); - listGroups(cb: (error: any, data: any) => any): void; - describeGroups( - consumerGroups: any, - cb: (error: any, data: any) => any, - ): void; - listTopics(cb: (error: any, data: any) => any): void; - createTopics( - topics: TopicConfigData[], - cb: (error: any, data: any) => any, - ): void; - describeConfigs( - payload: { resources: Resource[]; includeSynonyms: boolean }, - cb: (error: any, data: any) => any, - ): void; -} -export interface Resource { - resourceType: string; - resourceName: string; - configNames: string[]; -} - -export interface TopicConfigData { - topic: string; - partitions?: number; - replicationFactor?: number; - configEntry?: Array<{ name: string; value: string }>; -} -// # Interfaces - -export interface Message { - topic: string; - value: string | Buffer; - offset?: number; - partition?: number; - highWaterOffset?: number; - key?: string; -} - -export interface ProducerOptions { - requireAcks?: number; - ackTimeoutMs?: number; - partitionerType?: number; -} - -export interface KafkaClientOptions { - kafkaHost?: string; - connectTimeout?: number; - requestTimeout?: number; - autoConnect?: boolean; - connectRetryOptions?: RetryOptions; - sslOptions?: any; - clientId?: string; -} - -export interface RetryOptions { - retries?: number; - factor?: number; - minTimeout?: number; - maxTimeout?: number; - randomize?: boolean; -} - -export interface AckBatchOptions { - noAckBatchSize: number | null; - noAckBatchAge: number | null; -} - -export interface ZKOptions { - sessionTimeout?: number; - spinDelay?: number; - retries?: number; -} - -export interface ProduceRequest { - topic: string; - messages: any; // string[] | Array | string | KeyedMessage - key?: any; - partition?: number; - attributes?: number; -} - -export interface ConsumerOptions { - groupId?: string; - autoCommit?: boolean; - autoCommitIntervalMs?: number; - fetchMaxWaitMs?: number; - fetchMinBytes?: number; - fetchMaxBytes?: number; - fromOffset?: boolean; - encoding?: string; - keyEncoding?: string; -} - -export interface HighLevelConsumerOptions extends ConsumerOptions { - id?: string; - maxNumSegments?: number; - maxTickMessages?: number; - rebalanceRetry?: RetryOptions; -} - -export interface CustomPartitionAssignmentProtocol { - name: string; - version: number; - userData: {}; - assign(topicPattern: any, groupMembers: any, cb: (error: any, result: any) => void): void; -} - -export interface ConsumerGroupOptions { - kafkaHost?: string; - host?: string; - zk?: ZKOptions; - batch?: AckBatchOptions; - ssl?: boolean; - id?: string; - groupId: string; - sessionTimeout?: number; - protocol?: Array<"roundrobin" | "range" | CustomPartitionAssignmentProtocol>; - fromOffset?: "earliest" | "latest" | "none"; - outOfRangeOffset?: "earliest" | "latest" | "none"; - migrateHLC?: boolean; - migrateRolling?: boolean; - autoCommit?: boolean; - autoCommitIntervalMs?: number; - fetchMaxWaitMs?: number; - maxNumSegments?: number; - maxTickMessages?: number; - fetchMinBytes?: number; - fetchMaxBytes?: number; - retries?: number; - retryFactor?: number; - retryMinTimeout?: number; - connectOnReady?: boolean; -} - -export interface Topic { - topic: string; - offset?: number; - encoding?: string; - autoCommit?: boolean; -} - -export interface OffsetRequest { - topic: string; - partition?: number; - time?: number; - maxNum?: number; -} - -export interface OffsetCommitRequest { - topic: string; - partition?: number; - offset: number; - metadata?: string; -} - -export interface OffsetFetchRequest { - topic: string; - partition?: number; - offset?: number; -} - -export class TopicsNotExistError extends Error { - topics: string | string[]; -} - -export type CustomPartitioner = (partitions: number[], key: any) => number; diff --git a/types/kafka-node/kafka-node-tests.ts b/types/kafka-node/kafka-node-tests.ts deleted file mode 100644 index 7af6d11358..0000000000 --- a/types/kafka-node/kafka-node-tests.ts +++ /dev/null @@ -1,298 +0,0 @@ -import kafka = require('kafka-node'); - -const basicClient = new kafka.Client('localhost:2181/', 'sendMessage'); - -const optionsClient = new kafka.Client('localhost:2181/', 'sendMessage', { - sessionTimeout: 30000, - spinDelay: 1000, - retries: 0 -}, { - noAckBatchSize: 1000, - noAckBatchAge: 1000 * 10 -}, { - rejectUnauthorized: false -}); -optionsClient.topicExists(['topic'], (error: any) => { -}); -optionsClient.loadMetadataForTopics(['topic'], (error: any, data: any) => { -}); -optionsClient.refreshMetadata(['topic'], (error: any) => { -}); -optionsClient.close(); -optionsClient.sendOffsetCommitV2Request('group', 0, 'memberId', [], () => { -}); -optionsClient.close(() => { -}); - -const basicKafkaClient = new kafka.KafkaClient(); - -const optionsKafkaClient = new kafka.KafkaClient({ - kafkaHost: 'localhost:2181', - connectTimeout: 1000, - requestTimeout: 1000, - autoConnect: true, - sslOptions: {}, - clientId: "client id", - connectRetryOptions: { - retries: 5, - factor: 0, - minTimeout: 1000, - maxTimeout: 1000, - randomize: true - } -}); - -optionsKafkaClient.connect(); - -const optionsProducer = new kafka.Producer(basicClient, { - requireAcks: 0, - ackTimeoutMs: 0, - partitionerType: 0 -}); - -optionsKafkaClient.getListGroups((error: any, data: any) => { -}); -optionsKafkaClient.describeGroups([], (error: any, data: any) => { -}); - -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')] - }]; - - producer.send(messages, (err: Error) => { - }); - producer.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(); -}); - -const highLevelProducer = new kafka.HighLevelProducer(basicClient); -highLevelProducer.on('error', (error: Error) => { -}); -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(); -}); - -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.addTopics(['t1', 't2'], (err: any, added: any) => { -}); -consumer.addTopics([{topic: 't1', offset: 10}], (err: any, added: any) => { -}, true); - -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} -]); -consumer.resumeTopics([ - 'topic1', - {topic: 'topic2', partition: 0} -]); - -consumer.close(true, () => { -}); -consumer.close(() => { -}); - -const fetchRequests1 = [{topic: 'awesome'}]; -const hlConsumer = new kafka.HighLevelConsumer(basicClient, fetchRequests1, { - groupId: 'abcde', - autoCommit: true -}); - -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.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, () => { -}); -hlConsumer.close(() => { -}); - -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 -}; - -const consumerGroup = new kafka.ConsumerGroup(cgOptions, ['topic1']); -consumerGroup.on('error', (err) => { -}); -consumerGroup.on('message', (msg) => { -}); -consumerGroup.on('connect', () => { -}); -consumerGroup.close(true, () => { -}); - -const offset = new kafka.Offset(basicClient); - -offset.on('ready', () => { -}); - -offset.fetch([ - {topic: 't', partition: 0, time: Date.now(), maxNum: 1}, - {topic: 't'} -], (err: any, data: any) => { -}); - -offset.commit('groupId', [ - {topic: 't', partition: 0, offset: 10} -], (err, data) => { -}); - -offset.fetchCommits('groupId', [ - {topic: 't', partition: 0} -], (err, data) => { -}); - -offset.fetchLatestOffsets(['t'], (err, offsets) => { -}); -offset.fetchEarliestOffsets(['t'], (err, offsets) => { -}); - -const admin = new kafka.Admin(basicKafkaClient); -admin.listGroups((err, data) => { -}); -admin.describeGroups({}, (err, data) => { -}); -admin.createTopics([{ topic: 'testing' }], (err, data) => {}); -admin.listTopics((err, data) => { -}); - -const resource = { - resourceType: "", // 'broker' or 'topic' - resourceName: 'my-topic-name', - configNames: [] // specific config names, or empty array to return all, -}; -const payload = { - resources: [resource], - includeSynonyms: false // requires kafka 2.0+ -}; -admin.describeConfigs(payload, (err, res) => { - console.log(JSON.stringify(res, null, 1)); -}); diff --git a/types/kafka-node/tsconfig.json b/types/kafka-node/tsconfig.json deleted file mode 100644 index 154c35e191..0000000000 --- a/types/kafka-node/tsconfig.json +++ /dev/null @@ -1,23 +0,0 @@ -{ - "compilerOptions": { - "module": "commonjs", - "lib": [ - "es6" - ], - "noImplicitAny": true, - "noImplicitThis": true, - "strictNullChecks": true, - "strictFunctionTypes": true, - "baseUrl": "../", - "typeRoots": [ - "../" - ], - "types": [], - "noEmit": true, - "forceConsistentCasingInFileNames": true - }, - "files": [ - "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 deleted file mode 100644 index 75a24b3163..0000000000 --- a/types/kafka-node/tslint.json +++ /dev/null @@ -1,8 +0,0 @@ -{ - "extends": "dtslint/dt.json", - "rules": { - // TODOs - "no-any-union": false, - "no-unnecessary-class": false - } -}