Kafka HDFS Sink Connector with constant lag offset

Viewed 95

I have a Kafka HDFS Sink Connector with a constant offset lag. The kafka_consumergroup_group_lag from Kafka Lag Exporter is illustrated in the following figure:

Kafka connector offset lag

Note that the topic receives messages once a day, hence the spike. I would like the offset lag to go to zero, but as seen, the offset lag stabilizes at a value of ~833. How can I configure the connector to reach an offset lag of zero?

The connector configuration is given below

{
  "name": "my_connector",
  "config": {
    "connector.class": "io.confluent.connect.hdfs.HdfsSinkConnector",
    "tasks.max": "1",
    "retries": "2147483647",
    "topics": "my_kafka_topic",
    "format.class": "io.confluent.connect.hdfs.parquet.ParquetFormat",
    "partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner",
    "partition.duration.ms": "86400000",
    "path.format": "'date_id'=YYYYMMdd",
    "timezone": "UTC",
    "locale": "en-US",
    "timestamp.extractor": "RecordField",
    "timestamp.field": "message_timestamp",
    "compression.type": "snappy",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "errors.log.enable": "true",
    "errors.log.include.messages": "true",
    "errors.retry.delay.max.ms": "60000",
    "hadoop.conf.dir": "/var/run/configmaps/{{ stage }}",
    "hdfs.url": "{{ hdfs_url }}",
    "logs.dir": "{{ logs_dir }}",
    "topics.dir": "my_hdfs_path",
    "hdfs.authentication.kerberos": "true",
    "hdfs.namenode.principal": "{{ hdfs_namenode_principal }}",
    "connect.hdfs.principal": "{{ connect_hdfs_principal }}",
    "connect.hdfs.keytab": "{{ connect_hdfs_keytab }}",
    "flush.size": "600000",
    "rotate.interval.ms": "1600000",
    "transforms": "insertTS,formatTS",
    "transforms.insertTS.type": "org.apache.kafka.connect.transforms.InsertField$Value",
    "transforms.insertTS.timestamp.field": "message_timestamp",
    "transforms.formatTS.type": "org.apache.kafka.connect.transforms.TimestampConverter$Value",
    "transforms.formatTS.format": "yyyy-MM-dd'T'HH:mm:ss.SSSZ",
    "transforms.formatTS.field": "message_timestamp",
    "transforms.formatTS.target.type": "string"
  }
}

For topics receiving records more frequently, the connectors have no problem in receiving a zero offset lag (or close to zero):

Offset lag for my other connector

The configuration for this connector is identical:

{
  "name": "my_other_connector",
  "config": {
    "connector.class": "io.confluent.connect.hdfs.HdfsSinkConnector",
    "tasks.max": "1",
    "retries": "2147483647",
    "topics": "my_other_topic",
    "format.class": "io.confluent.connect.hdfs.parquet.ParquetFormat",
    "partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner",
    "partition.duration.ms": "86400000",
    "path.format": "'date_id'=YYYYMMdd",
    "timezone": "UTC",
    "locale": "en-US",
    "timestamp.extractor": "RecordField",
    "timestamp.field": "message_timestamp",
    "compression.type": "snappy",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "errors.log.enable": "true",
    "errors.log.include.messages": "true",
    "errors.retry.delay.max.ms": "60000",
    "hadoop.conf.dir": "/var/run/configmaps/{{ stage }}",
    "hdfs.url": "{{ hdfs_url }}",
    "logs.dir": "{{ logs_dir }}",
    "topics.dir": "my_other_hdfs_location",
    "hdfs.authentication.kerberos": "true",
    "hdfs.namenode.principal": "{{ hdfs_namenode_principal }}",
    "connect.hdfs.principal": "{{ connect_hdfs_principal }}",
    "connect.hdfs.keytab": "{{ connect_hdfs_keytab }}",
    "flush.size": "600000",
    "rotate.interval.ms": "1600000",
    "transforms": "insertTS,formatTS",
    "transforms.insertTS.type": "org.apache.kafka.connect.transforms.InsertField$Value",
    "transforms.insertTS.timestamp.field": "message_timestamp",
    "transforms.formatTS.type": "org.apache.kafka.connect.transforms.TimestampConverter$Value",
    "transforms.formatTS.format": "yyyy-MM-dd'T'HH:mm:ss.SSSZ",
    "transforms.formatTS.field": "message_timestamp",
    "transforms.formatTS.target.type": "string"
  }
}
0 Answers
Related