kafka-node: Provides its own types (#39580)

This commit is contained in:
Alexander T
2019-11-04 13:42:29 -08:00
committed by Nathan Shively-Sanders
parent 4cdec52083
commit 4c8ac68e40
5 changed files with 6 additions and 600 deletions
+6
View File
@@ -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",
-271
View File
@@ -1,271 +0,0 @@
// Type definitions for kafka-node 2.0
// Project: https://github.com/SOHU-Co/kafka-node/
// Definitions by: Daniel Imrie-Situnayake <https://github.com/dansitu>
// Bill <https://github.com/bkim54>
// Michael Haan <https://github.com/sfrooster>
// Amiram Korach <https://github.com/amiram>
// Insanehong <https://github.com/insanehong>
// Roger <https://github.com/rstpv>
// Definitions: https://github.com/DefinitelyTyped/DefinitelyTyped
/// <reference types="node" />
// # 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<OffsetFetchRequest | string>, 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<string|Topic> */): void;
resumeTopics(topics: any[] /* Array<string|Topic> */): 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<KeyedMessage> | 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;
-298
View File
@@ -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));
});
-23
View File
@@ -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"
]
}
-8
View File
@@ -1,8 +0,0 @@
{
"extends": "dtslint/dt.json",
"rules": {
// TODOs
"no-any-union": false,
"no-unnecessary-class": false
}
}