Commit 87fe461c authored by PrasadMoka's avatar PrasadMoka
Browse files

fix:LR-105 using kafka topic variables

Showing with 30 additions and 30 deletions
+30 -30
...@@ -83,7 +83,7 @@ base_config: | ...@@ -83,7 +83,7 @@ base_config: |
} }
} }
job { job {
env = "{{ env_name_bb }}" env = "{{ env_name }}"
enable.distributed.checkpointing = true enable.distributed.checkpointing = true
statebackend { statebackend {
blob { blob {
...@@ -141,11 +141,11 @@ activity-aggregate-updater: ...@@ -141,11 +141,11 @@ activity-aggregate-updater:
activity-aggregate-updater: |+ activity-aggregate-updater: |+
include file("/data/flink/conf/base-config.conf") include file("/data/flink/conf/base-config.conf")
kafka { kafka {
input.topic = {{ env_name_bb }}.coursebatch.job.request input.topic = {{ kafka_topic_course_batch_job_request }}
output.audit.topic = {{ env_name_bb }}.telemetry.raw output.audit.topic = {{ kafka_topic_telemetry_raw }}
output.failed.topic = {{ env_name_bb }}.activity.agg.failed output.failed.topic = {{ kafka_topic_activity_agg_failed }}
output.certissue.topic = {{ env_name_bb }}.issue.certificate.request output.certissue.topic = {{ kafka_topic_certificate_request }}
groupId = {{ env_name_bb }}-activity-aggregate-group groupId = {{ kafka_group_activity_agg }}
} }
task { task {
window.shards = {{ activity_agg_window_shards }} window.shards = {{ activity_agg_window_shards }}
...@@ -201,8 +201,8 @@ relation-cache-updater: ...@@ -201,8 +201,8 @@ relation-cache-updater:
relation-cache-updater: |+ relation-cache-updater: |+
include file("/data/flink/conf/base-config.conf") include file("/data/flink/conf/base-config.conf")
kafka { kafka {
input.topic = {{ env_name_bb }}.content.postpublish.request input.topic = {{ kafka_topic_content_publish_request }}
groupId = {{ env_name_bb }}-relation-cache-updater-group groupId = {{ kafka_group_relation_cache_updater }}
} }
task { task {
consumer.parallelism = {{ relation_cache_updater_consumer_parallelism }} consumer.parallelism = {{ relation_cache_updater_consumer_parallelism }}
...@@ -233,11 +233,11 @@ enrolment-reconciliation: ...@@ -233,11 +233,11 @@ enrolment-reconciliation:
enrolment-reconciliation: |+ enrolment-reconciliation: |+
include file("/data/flink/conf/base-config.conf") include file("/data/flink/conf/base-config.conf")
kafka { kafka {
input.topic = {{ env_name_bb }}.batch.enrolment.sync.request input.topic = {{ kafka_topic_enrolment_sync_request }}
output.audit.topic = {{ env_name_bb }}.telemetry.raw output.audit.topic = {{ kafka_topic_telemetry_raw }}
output.failed.topic = {{ env_name_bb }}.activity.agg.failed output.failed.topic = {{ kafka_topic_activity_agg_failed }}
output.certissue.topic = {{ env_name_bb }}.issue.certificate.request output.certissue.topic = {{ kafka_topic_certificate_request }}
groupId = {{ env_name_bb }}-enrolment-reconciliation-group groupId = {{ kafka_group_enrolment_reconciliation }}
} }
task { task {
restart-strategy.attempts = {{ restart_attempts }} # max 3 restart attempts restart-strategy.attempts = {{ restart_attempts }} # max 3 restart attempts
...@@ -279,10 +279,10 @@ collection-cert-pre-processor: ...@@ -279,10 +279,10 @@ collection-cert-pre-processor:
collection-cert-pre-processor: |+ collection-cert-pre-processor: |+
include file("/data/flink/conf/base-config.conf") include file("/data/flink/conf/base-config.conf")
kafka { kafka {
input.topic = {{ env_name_bb }}.issue.certificate.request input.topic = {{ kafka_topic_certificate_request }}
output.topic = {{ env_name_bb }}.generate.certificate.request output.topic = {{ kafka_topic_generate_certificate_request }}
output.failed.topic = {{ env_name_bb }}.issue.certificate.failed output.failed.topic = {{ kafka_topic_certificate_failed }}
groupId = {{ env_name_bb }}-collection-cert-pre-processor-group groupId = {{ kafka_group_collection_pre_processor }}
} }
task { task {
restart-strategy.attempts = {{ restart_attempts }} # max 3 restart attempts restart-strategy.attempts = {{ restart_attempts }} # max 3 restart attempts
...@@ -329,9 +329,9 @@ collection-certificate-generator: ...@@ -329,9 +329,9 @@ collection-certificate-generator:
collection-certificate-generator: |+ collection-certificate-generator: |+
include file("/data/flink/conf/base-config.conf") include file("/data/flink/conf/base-config.conf")
kafka { kafka {
input.topic = {{ env_name_bb }}.generate.certificate.request input.topic = {{ kafka_topic_generate_certificate_request }}
output.audit.topic = {{ env_name_bb }}.telemetry.raw output.audit.topic = {{ kafka_topic_telemetry_raw }}
groupId = {{ env_name_bb }}-certificate-generator-group groupId = {{ kafka_group_certificate_generator }}
} }
task { task {
restart-strategy.attempts = {{ restart_attempts }} # max 3 restart attempts restart-strategy.attempts = {{ restart_attempts }} # max 3 restart attempts
...@@ -375,10 +375,10 @@ merge-user-courses: ...@@ -375,10 +375,10 @@ merge-user-courses:
merge-user-courses: |+ merge-user-courses: |+
include file("/data/flink/conf/base-config.conf") include file("/data/flink/conf/base-config.conf")
kafka { kafka {
input.topic = {{ env_name_bb }}.lms.user.account.merge input.topic = {{ kafka_topic_lms_user_account }}
output.failed.topic = {{ env_name_bb }}.learning.events.failed output.failed.topic = {{ kafka_topic_learning_failed }}
groupId = {{ env_name_bb }}-merge-courses-group groupId = {{ kafka_group_merge_courses }}
output.course.batch.updater.topic = {{ env_name_bb }}.coursebatch.job.request output.course.batch.updater.topic = {{ kafka_topic_course_batch_job_request }}
} }
task { task {
consumer.parallelism = {{ merge_user_courses_consumer_parallelism }} consumer.parallelism = {{ merge_user_courses_consumer_parallelism }}
...@@ -408,10 +408,10 @@ assessment-aggregator: ...@@ -408,10 +408,10 @@ assessment-aggregator:
producer.broker-servers = "{{ kafka_brokers }}" producer.broker-servers = "{{ kafka_brokers }}"
consumer.broker-servers = "{{ kafka_brokers }}" consumer.broker-servers = "{{ kafka_brokers }}"
zookeeper = "{{ zookeepers }}" zookeeper = "{{ zookeepers }}"
input.topic = {{ env_name_bb }}.telemetry.assess input.topic = {{ kafka_topic_assessment }}
failed.topic= {{ env_name_bb }}.telemetry.assess.failed failed.topic= {{ kafka_topic_assessment_failed }}
groupId = {{ env_name_bb }}-assessment-aggregator-group groupId = {{ kafka_group_assessment_aggregator }}
output.certissue.topic = {{ env_name_bb }}.issue.certificate.request output.certissue.topic = {{ kafka_topic_certificate_request }}
} }
task { task {
consumer.parallelism = {{ assessaggregator_consumer_parallelism }} consumer.parallelism = {{ assessaggregator_consumer_parallelism }}
...@@ -450,8 +450,8 @@ notification-job: ...@@ -450,8 +450,8 @@ notification-job:
notification-job: |+ notification-job: |+
include file("/data/flink/conf/base-config.conf") include file("/data/flink/conf/base-config.conf")
kafka { kafka {
input.topic = {{ env_name_bb }}.lms.notification input.topic = {{ kafka_topic_lms_notification }}
groupId = {{ env_name_bb }}-lms.notification groupId = {{ kafka_group_lms_notification }}
} }
task { task {
restart-strategy.attempts = {{ restart_attempts }} # max 3 restart attempts restart-strategy.attempts = {{ restart_attempts }} # max 3 restart attempts
......
Supports Markdown
0% or .
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment