Merge "[MON] Implements multithreading for message consumption"
diff --git a/MANIFEST.in b/MANIFEST.in
index 0887e73..485fc0a 100644
--- a/MANIFEST.in
+++ b/MANIFEST.in
@@ -21,6 +21,6 @@
include requirements.txt
include README.rst
-recursive-include osm_mon *.py
+recursive-include osm_mon *.py *.xml
recursive-include devops-stages *
recursive-include test *.py
diff --git a/policy_module/osm_policy_module/core/agent.py b/policy_module/osm_policy_module/core/agent.py
index 3a3b20c..cdd5dfc 100644
--- a/policy_module/osm_policy_module/core/agent.py
+++ b/policy_module/osm_policy_module/core/agent.py
@@ -45,11 +45,10 @@
# Initialize Kafka consumer
log.info("Connecting to Kafka server at %s", kafka_server)
- # TODO: Add logic to handle deduplication of messages when using group_id.
- # See: https://stackoverflow.com/a/29836412
consumer = KafkaConsumer(bootstrap_servers=kafka_server,
key_deserializer=bytes.decode,
- value_deserializer=bytes.decode)
+ value_deserializer=bytes.decode,
+ group_id="pm-consumer")
consumer.subscribe(['lcm_pm', 'alarm_response'])
for message in consumer:
diff --git a/setup.py b/setup.py
index 20d4068..648d9e9 100644
--- a/setup.py
+++ b/setup.py
@@ -50,9 +50,6 @@
license=_license,
packages=[_name],
package_dir={_name: _name},
- package_data={_name: ['osm_mon/core/message_bus/*.py', 'osm_mon/core/models/*.json',
- 'osm_mon/plugins/OpenStack/Aodh/*.py', 'osm_mon/plugins/OpenStack/Gnocchi/*.py',
- 'osm_mon/plugins/vRealiseOps/*', 'osm_mon/plugins/CloudWatch/*']},
scripts=['osm_mon/plugins/vRealiseOps/vROPs_Webservice/vrops_webservice',
'osm_mon/core/message_bus/common_consumer.py'],
install_requires=parse_requirements('requirements.txt'),