标签: confluent-platform

Kafka JDBC 连接器加载所有数据,然后增量

我试图弄清楚如何最初从查询中获取所有数据,然后仅使用 kafka 连接器增量更改。这样做的原因是我想将所有数据加载到弹性搜索中,然后使 es 与我的 kafka 流同步。目前,我首先使用带有模式 = 批量的连接器来执行此操作,然后将其更改为时间戳。这工作正常。

但是,如果我们想将所有数据重新加载到 Streams 和 ES,这意味着我们必须编写一些脚本来以某种方式清理或删除 kafka 流和 es 索引数据,修改连接 ini 以将模式设置为批量,重新启动所有内容,给出是时候加载所有数据,然后再次将脚本修改为时间戳模式,然后再次重新启动所有内容(需要这样一个脚本的原因是偶尔,批量更新会通过我们尚无法控制的 etl 过程来纠正历史数据,并且此过程不会更新时间戳)

有没有人做类似的事情并找到了更优雅的解决方案?

elasticsearch apache-kafka apache-kafka-connect confluent-platform

7
推荐指数
1
解决办法
1949
查看次数

如何在 Confluence kafka C# 库中获取 Kafka 主题的最新偏移量?

我正在使用 Confluence kafka C# 客户端。如何获取此主题中消耗的最新偏移量?

c# apache-kafka kafka-consumer-api confluent-platform

7
推荐指数
2
解决办法
1万
查看次数

Kafka SMT ValueToKey - 如何使用多个值作为键?

我正在使用 Confluent JDBCSourceConnector 从 Oracle 表中读取数据。我正在尝试使用 SMT 生成由 3 个连接字段组成的密钥。

transforms=createKey
transforms.createKey.type=org.apache.kafka.connect.transforms.ValueToKey
transforms.createKey.fields=BUS_OFC_ID_CD,SO_TYPE,SO_NO
Run Code Online (Sandbox Code Playgroud)

使用上面的转换,我得到了这样的东西:

{"BUS_OFC_ID_CD":"111","SO_TYPE":"I","SO_NO":"55555"}
Run Code Online (Sandbox Code Playgroud)

我想要类似的东西:

111I55555
Run Code Online (Sandbox Code Playgroud)

关于如何仅连接值的任何想法?

oracle jdbc apache-kafka apache-kafka-connect confluent-platform

7
推荐指数
1
解决办法
2346
查看次数

如何使用 Spring-Kafka 通过 Confluent Schema 注册表读取 AVRO 消息?

如何使用 Spring-Kafka 通过 Confluent Schema 注册表读取 AVRO 消息?有样品吗?我在官方参考文件中找不到它。

avro spring-kafka confluent-schema-registry confluent-platform

7
推荐指数
2
解决办法
1万
查看次数

如何使用 Kafka Connect JDBC 来获取具有多个包含同名表的模式的 PostgreSQL?

我需要从具有约 2000 个模式的 PostgreSQL 数据库中获取数据。所有模式都包含相同的表(它是一个多租户应用程序)。

连接器配置如下:

{
  "name": "postgres-source",
  "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
  "timestamp.column.name": "updated",
  "incrementing.column.name": "id",
  "connection.password": "********",
  "tasks.max": "1",
  "mode": "timestamp+incrementing",
  "topic.prefix": "postgres-source-",
  "connection.user": "*********",
  "poll.interval.ms": "3600000",
  "numeric.mapping": "best_fit",
  "connection.url": "jdbc:postgresql://*******:5432/*******",
  "table.whitelist": "table1",
  "value.converter": "org.apache.kafka.connect.json.JsonConverter",
  "value.converter.schemas.enable": "false",
  "key.converter": "org.apache.kafka.connect.json.JsonConverter",
  "key.converter.schemas.enable":"false"
}
Run Code Online (Sandbox Code Playgroud)

使用此配置,我收到此错误:

“连接器使用非限定表名作为主题名,并检测到重复的非限定表名。这可能导致主题中的混合数据类型和下游处理错误。为防止此类处理错误,JDBC Source 连接器在执行时无法启动检测重复的表名配置”

显然,连接器不想将多个同名表中的数据发布到单个主题。

这对我来说无关紧要,它可以转到单个主题或多个主题(每个模式一个)。

作为附加信息,如果我添加:

"schema.pattern": "schema1" 
Run Code Online (Sandbox Code Playgroud)

到配置,连接器工作并且来自指定模式和表的数据被复制。

有没有办法复制包含同名表的多个模式?

谢谢

postgresql jdbc apache-kafka apache-kafka-connect confluent-platform

7
推荐指数
0
解决办法
1216
查看次数

如何在 C# 中为泛型类型创建 Avro 模式?

我正在使用 Kafka 和Confluent 的 .NET 客户端开发 .NET Standard pub/sub 包。我的制作人有以下界面。

IEventPublisher.cs

public interface IEventPublisher<T>
{
    bool Publish(Event<T> evnt);
}
Run Code Online (Sandbox Code Playgroud)

我的KafkaEventPublisher<T>类实现了这个接口,并且正在发布的有效负载 T 被包裹在一个Event<T>信封中。

事件.cs

public class Event<T>
{
    // Some other properties 

    public T Payload { get; set; }
}
Run Code Online (Sandbox Code Playgroud)

我的组件的初始实现不使用 Avro 序列化程序或架构注册表。它Event<T>使用 JSON序列化并将Newtonsoft.Json字符串生成到 Kafka 主题。这样做的好处是有效载荷实际上可以是任何对象。主题名称是对象的完全限定类名称,因此主题保证是同类的。缺点是有效载荷对 Kafka 是不透明的。

我现在正致力于从Newtonsoft.JsonAvro 和架构注册表转换。这似乎表明我的模型不能再是任何东西了。它们必须专门编写以通过实现ISpecificRecord接口来允许 Avro 序列化。如果这是真的,这并不理想,但我可以接受。

我似乎无法弄清楚的问题是如何将Event<T>信封合并到 Avro 模式中。有没有办法将模式嵌套在另一个模式中?我所有的具体模式都应该定义Event<T>信封吗?GenericRecord改用并尝试在运行时使用反射自动生成模式是否明智,这样我的组件就可以重新接受任何对象作为有效负载?

.net c# avro apache-kafka confluent-platform

7
推荐指数
1
解决办法
677
查看次数

当官方控制台消费者正常工作时,汇合的kafka python Consumer.poll() 总是返回 None

当设置超时调用时,confluence-kafka python 客户端的实例Consumer始终返回 None 。poll()

该主题确实包含一些消息,并且官方控制台消费者工作正常:

$ vim ~/client.properties

security.protocol=SASL_PLAINTEXT
sasl.mechanism=SCRAM-SHA-512
sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required username=XXXXXXXXXX password="XXXXXXXXXX";

$ ~/kafka_2.13-2.4.0/bin/kafka-console-consumer.sh --topic my_topic --bootstrap-server somehost.:30742 --from-beginning --consumer.config ~/client.properties --group somenewgroup
msg1
msg2
msg3
Run Code Online (Sandbox Code Playgroud)

Consumer.poll()方法总是返回 None。即使当我将密码或主机更改为无效值时,它也会返回 None 。

我为消费者设置了一个记录器,但没有记录任何内容。

python代码如下

consumer=Consumer({'sasl.mechanisms': "SCRAM-SHA-512",
                   'security.protocol': 'SASL_PLAINTEXT',
                   'sasl.username': 'XXXXXXXXXX',
                   'sasl.password': 'XXXXXXXXXX',
                   'bootstrap.servers': 'somehost.:30742',
                   "group.id":"somenewgroup",
                   'auto.offset.reset': 'beginning',
                   'logger':logger
                   },logger=logger)

consumer.subscribe(["my_topic"])

while True:
    msg = consumer.poll(timeout=1.0)
    print("poll success")
    if msg is None:print("msg is None!")
Run Code Online (Sandbox Code Playgroud)

python apache-kafka kafka-consumer-api confluent-platform

7
推荐指数
0
解决办法
1348
查看次数

Kafka 连接,Bootstrap 代理断开连接


我试图设置 Kafka Connect 以运行 ElasticsearchSinkConnector。
Kafka 设置,由 3 个使用 Kerberos、SSL 和 ACL 保护的代理组成。

到目前为止,我一直在尝试使用 docker/docker-compose(Confluent docker-image 5.4 with Kafka 2.4)连接到远程 kafka 安装(Kafka 2.0.1 - 实际上是我们的生产环境)运行连接框架和 elasticserch-server 本地)。

KAFKA_OPTS: -Djava.security.krb5.conf=/etc/kafka-connect/secrets/krb5.conf
      CONNECT_BOOTSTRAP_SERVERS: srv-kafka-1.XXX.com:9093,srv-kafka-2.XXX.com:9093,srv-kafka-3.XXX.com:9093
      CONNECT_REST_ADVERTISED_HOST_NAME: kafka-connect
      CONNECT_REST_PORT: 8083
      CONNECT_GROUP_ID: user-grp
      CONNECT_CONFIG_STORAGE_TOPIC: test.internal.connect.configs
      CONNECT_OFFSET_STORAGE_TOPIC: test.internal.connect.offsets
      CONNECT_STATUS_STORAGE_TOPIC: test.internal.connect.status
      CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR: 1
      CONNECT_STATUS_STORAGE_REPLICATION_FACTOR: 1
      CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR: 1
      CONNECT_KEY_CONVERTER: org.apache.kafka.connect.json.JsonConverter
      CONNECT_VALUE_CONVERTER: org.apache.kafka.connect.json.JsonConverter
      CONNECT_INTERNAL_KEY_CONVERTER: org.apache.kafka.connect.json.JsonConverter
      CONNECT_INTERNAL_VALUE_CONVERTER: org.apache.kafka.connect.json.JsonConverter
      CONNECT_ZOOKEEPER_CONNECT: srv-kafka-1.XXX.com:2181,srv-kafka-2.XXX.com:2181,srv-kafka-3.XXX.com:2181
      CONNECT_SECURITY_PROTOCOL: SASL_SSL
      CONNECT_SASL_KERBEROS_SERVICE_NAME: "kafka"
      CONNECT_SASL_JAAS_CONFIG: com.sun.security.auth.module.Krb5LoginModule required \
                                useKeyTab=true \
                                storeKey=true \
                                keyTab="/etc/kafka-connect/secrets/kafka-connect.keytab" \
                                principal="<principal>;
      CONNECT_SASL_MECHANISM: GSSAPI
      CONNECT_SSL_TRUSTSTORE_LOCATION: <path_to_truststore.jks>
      CONNECT_SSL_TRUSTSTORE_PASSWORD: <PWD> …
Run Code Online (Sandbox Code Playgroud)

apache-kafka apache-kafka-connect confluent-platform

7
推荐指数
1
解决办法
1万
查看次数

Apple M1:无法使用docker启动Confluence控制中心

我想按照文档中的规定使用docker启动Confluence控制中心:

https://docs.confluence.io/platform/current/quickstart/ce-docker-quickstart.html

这是docker-compose.yaml他们提供的文件;

---
version: '2'
services:
  zookeeper:
    image: confluentinc/cp-zookeeper:5.5.0
    hostname: zookeeper
    container_name: zookeeper
    ports:
      - "2181:2181"
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000

  broker:
    image: confluentinc/cp-server:5.5.0
    hostname: broker
    container_name: broker
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181'
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://broker:29092,PLAINTEXT_HOST://localhost:9092
      KAFKA_METRIC_REPORTERS: io.confluent.metrics.reporter.ConfluentMetricsReporter
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
      KAFKA_CONFLUENT_LICENSE_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      CONFLUENT_METRICS_REPORTER_BOOTSTRAP_SERVERS: broker:29092
      CONFLUENT_METRICS_REPORTER_ZOOKEEPER_CONNECT: zookeeper:2181
      CONFLUENT_METRICS_REPORTER_TOPIC_REPLICAS: 1
      CONFLUENT_METRICS_ENABLE: 'true'
      CONFLUENT_SUPPORT_CUSTOMER_ID: 'anonymous'

  schema-registry:
    image: confluentinc/cp-schema-registry:5.5.0
    hostname: schema-registry
    container_name: schema-registry
    depends_on:
      - zookeeper
      - broker …
Run Code Online (Sandbox Code Playgroud)

apache-kafka docker confluent-control-center confluent-platform apple-m1

7
推荐指数
1
解决办法
7765
查看次数

Kafka Connect:没有为连接器创建任务

我们正在使用 Debezium (MongoDB) 和 Confluent S3 连接器以分布式模式运行 Kafka Connect(Confluent Platform 5.4,即 Kafka 2.4)。通过 REST API 添加新连接器时,连接器创建为 RUNNING 状态,但不会为连接器创建任何任务。

暂停和恢复连接器无济于事。当我们停止所有工人然后再次启动它们时,任务被创建,一切都按预期运行。

该问题不是由连接器插件引起的,因为我们看到 Debezium 和 S3 连接器的行为相同。同样在调试日志中,我可以看到 Debezium 正确地从 Connector.taskConfigs() 方法返回任务配置。

有人可以告诉我该怎么做,我们可以在不重新启动工作人员的情况下添加连接器吗?谢谢。

配置详情

集群有 3 个节点,具有以下connect-distributed.properties

bootstrap.servers=kafka-broker-001:9092,kafka-broker-002:9092,kafka-broker-003:9092,kafka-broker-004:9092
group.id=tdp-QA-connect-cluster

key.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=false
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false

internal.key.converter=org.apache.kafka.connect.json.JsonConverter
internal.value.converter=org.apache.kafka.connect.json.JsonConverter
internal.key.converter.schemas.enable=false
internal.value.converter.schemas.enable=false

offset.storage.topic=connect-offsets-qa
offset.storage.replication.factor=3
offset.storage.partitions=5

config.storage.topic=connect-configs-qa
config.storage.replication.factor=3

status.storage.topic=connect-status-qa
status.storage.replication.factor=3
status.storage.partitions=3

offset.flush.interval.ms=10000

rest.host.name=tdp-QA-kafka-connect-001
rest.port=10083
rest.advertised.host.name=tdp-QA-kafka-connect-001
rest.advertised.port=10083

plugin.path=/opt/kafka-connect/plugins,/usr/share/java/

security.protocol=SSL
ssl.truststore.location=/etc/kafka/ssl/kafka-connect.truststore.jks
ssl.truststore.password=<secret>
ssl.endpoint.identification.algorithm=
producer.security.protocol=SSL
producer.ssl.truststore.location=/etc/kafka/ssl/kafka-connect.truststore.jks
producer.ssl.truststore.password=<secret>
consumer.security.protocol=SSL
consumer.ssl.truststore.location=/etc/kafka/ssl/kafka-connect.truststore.jks
consumer.ssl.truststore.password=<secret>

max.request.size=20000000
max.partition.fetch.bytes=20000000
Run Code Online (Sandbox Code Playgroud)

连接器配置

Debezium 示例:

{
  "name": "qa-mongodb-comp-converter-task|1",
  "config": {
    "connector.class": …
Run Code Online (Sandbox Code Playgroud)

apache-kafka apache-kafka-connect confluent-platform

6
推荐指数
1
解决办法
1653
查看次数