我试图弄清楚如何最初从查询中获取所有数据,然后仅使用 kafka 连接器增量更改。这样做的原因是我想将所有数据加载到弹性搜索中,然后使 es 与我的 kafka 流同步。目前,我首先使用带有模式 = 批量的连接器来执行此操作,然后将其更改为时间戳。这工作正常。
但是,如果我们想将所有数据重新加载到 Streams 和 ES,这意味着我们必须编写一些脚本来以某种方式清理或删除 kafka 流和 es 索引数据,修改连接 ini 以将模式设置为批量,重新启动所有内容,给出是时候加载所有数据,然后再次将脚本修改为时间戳模式,然后再次重新启动所有内容(需要这样一个脚本的原因是偶尔,批量更新会通过我们尚无法控制的 etl 过程来纠正历史数据,并且此过程不会更新时间戳)
有没有人做类似的事情并找到了更优雅的解决方案?
elasticsearch apache-kafka apache-kafka-connect confluent-platform
我正在使用 Confluence kafka C# 客户端。如何获取此主题中消耗的最新偏移量?
我正在使用 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
如何使用 Spring-Kafka 通过 Confluent Schema 注册表读取 AVRO 消息?有样品吗?我在官方参考文件中找不到它。
avro spring-kafka confluent-schema-registry confluent-platform
我需要从具有约 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
我正在使用 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改用并尝试在运行时使用反射自动生成模式是否明智,这样我的组件就可以重新接受任何对象作为有效负载?
当设置超时调用时,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)
我试图设置 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) 我想按照文档中的规定使用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
我们正在使用 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 ×9
avro ×2
c# ×2
jdbc ×2
.net ×1
apple-m1 ×1
docker ×1
oracle ×1
postgresql ×1
python ×1
spring-kafka ×1