需要检查生产中监控Kafka的工具.此外,工具不需要许可证或重型硬件.特别是我需要一个工具来评估消费者对主题的偏差,主题的健康状况.
Oracle SQL 支持START WITH表达式。例如,
CREATE VIEW customers AS
SELECT LEVEL lvl, customer_code, customer_desc, customer_category
FROM customers_master
START WITH some_column = '0'
CONNECT BY PRIOR CUSTOMER_CODE = PARENT_CUSTOMER_CODE;
Run Code Online (Sandbox Code Playgroud)
如果表包含分层数据,则可以使用分层查询子句按分层顺序选择行。
START WITH指定层次结构的根行。
CONNECT BY指定层次结构的父行和子行之间的关系。
MS-SQL 是否有等效的表达式?
生产者通过设置 Kafka Broker 列表来发送消息,如下所示。
props.put("bootstrap.servers", "127.0.0.1:9092,127.0.0.1:9092,127.0.0.1:9092");
Run Code Online (Sandbox Code Playgroud)
我想知道“生产者”如何知道三个经纪人中的哪一个知道哪个有分区领导者。对于典型的分布式服务器,要么你有一个承载服务器,要么有一个虚拟IP,但对于Kafka,它是如何加载的?生产者程序是否尝试随机连接到一个代理并寻找具有分区领导者的代理?
根据文档,Old Consumer Configsconsumer.id默认包含一个为 null的集合:
consumer.id
默认值:null
描述:如果未设置则自动生成。
是否可以consumer.id为New Kafka Consumer 设置 ,如果可以,我该如何实现?
假设我要从 kafka 生产者向 kafka 消费者发送一些消息,那么它将存储在哪里?是否有用于存储消息的数据库?消息存储多长时间?
任何人都可以请解释一下。
我想使用Confluent的REST API运行JDBC源连接器.虽然独立模式使用以下属性文件可以完美运行:
name=source-mysql-test
connector.class=io.confluent.connect.jdbc.JdbcSourceConnector
tasks.max=1
connection.url=jdbc:mysql://localhost:3306/kafka
connection.user=myuser
connection.password=mypass
table.whitelist=MY_TABLE
# Pull all rows based on timestamp
mode=timestamp
timestamp.column.name=ROWVERSION
validate.non.null=false
# The Kafka topic will be made up of this prefix, plus the table name.
topic.prefix=MYSQL-
table.types=TABLE,VIEW
poll.interval.ms=1000
Run Code Online (Sandbox Code Playgroud)
我无法使用REST API运行连接器.这是电话:
curl -X POST -H "Content-Type: application/json" --data '{"name": "source-mysql-test", "config": {"connector.class":"io.confluent.connect.jdbc.JdbcSourceConnector", "tasks.max":"1", "connection.url":"jdbc:mysql://localhost:3306/kafka","connection.user":"myuser","connection.password":"mypass", "table.whitelist":"MY_TABLE", "mode":"timestamp", "timestamp.column.name":"ROWVERSION", "validate.non.null":"false", "topic.prefix":"MYSQL-", "table.types":"TABLE,VIEW", "poll.interval.ms":"1000" }}' http://localhost:8083/connectors
Run Code Online (Sandbox Code Playgroud)
以下是回复:
{
"error_code": 400,
"message": "Connector configuration is invalid and contains the following 2 error(s):\nInvalid value com.mysql.jdbc.exceptions.jdbc4.MySQLNonTransientConnectionException: Could not …Run Code Online (Sandbox Code Playgroud) 我正在通过kafka connect,我正在尝试获取这些概念.
让我们说我有kafka集群(节点k1,k2和k3)设置并且它正在运行,现在我想在不同的节点中运行kafka connect worker,比如分布式模式下的c1和c2.
几个问题.
1)要在分布式模式下运行或启动kafka connect,我需要使用../bin/connect-distributed.shkakfa集群节点中可用的命令,所以我需要从任何一个kafka集群节点启动kafka connect?或者我启动kafka connect的任何节点都需要有kafka二进制文件才能使用../bin/connect-distributed.sh
2)我需要将我的连接器插件复制到我执行第1步的任何kafka集群节点(或所有集群节点?)?
3)在工作节点上启动jvm进程之前,kafka如何将这些连接器插件复制到工作节点?因为插件是具有我的任务代码的插件,需要将其复制到worker才能在worker中启动进程.
4)我是否需要在连接群集节点c1和c2中安装任何东西,比如需要安装java或任何相关的kafka连接?
5)在某些地方,它说使用汇合平台,但我想先用apache kafka connect启动它.
有些人可以通过一些亮点甚至指向一些资源也会有所帮助.
谢谢.
我想用 python 获取 Kafka 消费者组列表,但我不能。
我使用zookeeper python客户端(kazoo),但消费者组列表为空,因为这种方法适用于旧消费者,而我们没有使用旧消费者。
如何使用python代码获取消费者组列表?
./kafka-consumer-groups.sh -bootstrap-server localhost:9092 -list
Run Code Online (Sandbox Code Playgroud) kafka 领导者是自己分区还是经纪人?我最初的理解是,它们是充当读/写代理的分区,然后将它们的价值交给 ISR。
但是最近我听到他们提到他们好像发生在“经纪人”级别,因此我很困惑。
我知道还有其他帖子旨在回答这个问题,但那里的答案没有帮助。
我正在运行 kafka-connect 分布式设置。
我正在使用单机/进程设置(仍处于分布式模式)进行测试,效果很好,现在我正在使用 3 个节点(和 3 个连接进程),日志不包含错误,但是当我提交 s3-connector 时通过rest-api请求,它返回:{"error_code":409,"message":"Cannot complete request because of a conflicting operation (e.g. worker rebalance)"}。
当我停止其中一个节点上的 kafka-connect 进程时,我实际上可以提交作业并且一切正常。
我的集群中有 3 个代理,主题的分区号是 32。
这是我尝试启动的连接器:
{
"name": "s3-sink-new-2",
"config": {
"connector.class": "io.confluent.connect.s3.S3SinkConnector",
"tasks.max": "32",
"topics": "rawEventsWithoutAttribution5",
"s3.region": "us-east-1",
"s3.bucket.name": "dy-raw-collection",
"s3.part.size": "64000000",
"flush.size": "10000",
"storage.class": "io.confluent.connect.s3.storage.S3Storage",
"format.class": "io.confluent.connect.s3.format.avro.AvroFormat",
"schema.generator.class": "io.confluent.connect.storage.hive.schema.DefaultSchemaGenerator",
"partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner",
"partition.duration.ms": "60000",
"path.format": "\'year\'=YYYY/\'month\'=MM/\'day\'=dd/\'hour\'=HH",
"locale": "US",
"timezone": "GMT",
"timestamp.extractor": "RecordField",
"timestamp.field": "procTimestamp",
"name": "s3-sink-new-2"
}
}
Run Code Online (Sandbox Code Playgroud)
日志中没有任何内容表明有问题,我真的迷失在这里。
apache-kafka ×9
confluent ×2
broker ×1
connect-by ×1
jdbc ×1
kafka-topic ×1
leader ×1
monitoring ×1
mysql ×1
oracle ×1
python ×1
sql-server ×1
t-sql ×1