import json import os import shelve from pprint import pprint import asyncio import jks import ssl from aiokafka import AIOKafkaProducer from aiosfstream import Client from aiosfstream import RefreshTokenAuthenticator BOOTSTRAP_SERVERS = os.environ.get('BOOSTSTRAP_SERVERS') TOPIC = os.environ.get('KAFKA_TOPIC', 'event.aminakovTest') KESTORE_PASS = os.environ.get('JKS_KESTORE_PASSWORD', 'chageit') JKS_PATH = os.environ.get('JKS_KEYSTORE_PATH') async def stream_events(): keystore = jks.KeyStore.load(JKS_PATH, KESTORE_PASS) ctx = ssl.SSLContext(protocol=ssl.PROTOCOL_TLSv1_2) for alias, cert in keystore.certs.items(): ctx.load_verify_locations(cadata=cert.cert) producer = AIOKafkaProducer( bootstrap_servers=BOOTSTRAP_SERVERS, security_protocol='SSL', ssl_context=ctx) await producer.start() with shelve.open('replay.db') as replay: auth = RefreshTokenAuthenticator( consumer_key=os.environ.get('SALESFORCE_CONSUMER_KEY'), consumer_secret=os.environ.get('SALESFORCE_CONSUMER_SECRET'), refresh_token=os.environ.get('SALESFORCE_REFRESH_TOKEN'), sandbox=False) async with Client(authenticator=auth, replay=replay) as client: # subscribe to the platform event using CometD await client.subscribe('/data/ContactChangeEvent') # listen for incoming messages async for message in client: topic = message['channel'] data = message['data']['payload'] print(f'{topic}') pprint(data) await producer.send(TOPIC, json.dumps(data).encode()) if __name__ == '__main__': loop = asyncio.get_event_loop() loop.run_until_complete(stream_events())