{"id":47087,"date":"2023-05-21T12:22:22","date_gmt":"2022-12-29T02:13:02","guid":{"rendered":"https:\/\/www.silicloud.com\/zh\/blog\/2022%e5%b9%b410%e6%9c%8819%e6%97%a5%e7%9a%84cloudformation-codepipeline%e5%a4%87%e5%bf%98%e5%bd%95%e3%80%82\/"},"modified":"2024-01-15T12:16:44","modified_gmt":"2024-01-15T04:16:44","slug":"2022%e5%b9%b410%e6%9c%8819%e6%97%a5%e7%9a%84cloudformation-codepipeline%e5%a4%87%e5%bf%98%e5%bd%95%e3%80%82","status":"publish","type":"post","link":"https:\/\/www.silicloud.com\/zh\/blog\/2022%e5%b9%b410%e6%9c%8819%e6%97%a5%e7%9a%84cloudformation-codepipeline%e5%a4%87%e5%bf%98%e5%bd%95%e3%80%82\/","title":{"rendered":"2022\u5e7410\u670819\u65e5\u7684CloudFormation CodePipeline\u5907\u5fd8\u5f55"},"content":{"rendered":"<pre class=\"post-pre\"><code><span class=\"o\">---<\/span>\r\n<span class=\"c1\">#!\/bin\/env python<\/span>\r\n\r\n<span class=\"n\">import<\/span> <span class=\"n\">json<\/span>\r\n<span class=\"n\">import<\/span> <span class=\"n\">logging<\/span>\r\n<span class=\"n\">import<\/span> <span class=\"n\">os<\/span>\r\n<span class=\"n\">import<\/span> <span class=\"n\">boto3<\/span>\r\n<span class=\"n\">import<\/span> <span class=\"n\">datetime<\/span>\r\n<span class=\"n\">from<\/span> <span class=\"n\">bson<\/span> <span class=\"n\">import<\/span> <span class=\"n\">json_util<\/span>\r\n<span class=\"n\">from<\/span> <span class=\"n\">pymongo<\/span> <span class=\"n\">import<\/span> <span class=\"no\">MongoClient<\/span>\r\n<span class=\"n\">from<\/span> <span class=\"n\">pymongo<\/span><span class=\"p\">.<\/span><span class=\"nf\">errors<\/span> <span class=\"n\">import<\/span> <span class=\"no\">OperationFailure<\/span>\r\n<span class=\"n\">import<\/span> <span class=\"n\">urllib<\/span><span class=\"p\">.<\/span><span class=\"nf\">request<\/span>                                                    \r\n\r\n<span class=\"s2\">\"\"\"\r\nRead data from a DocumentDB collection's change stream and replicate that data to MSK.\r\n\r\nRequired environment variables:\r\nDOCUMENTDB_URI: The URI of the DocumentDB cluster to stream from.\r\nDOCUMENTDB_SECRET: Secret Name of the credentials for the DocumentDB cluster in Secrets Manager\r\nSTATE_COLLECTION: The name of the collection in which to store sync state.\r\nSTATE_DB: The name of the database in which to store sync state.\r\nWATCHED_COLLECTION_NAME: The name of the collection to watch for changes.\r\nWATCHED_DB_NAME: The name of the database to watch for changes.\r\nIterations_per_sync: How many events to process before syncing state.\r\nDocuments_per_run: The max for the iterator loop. \r\nSNS_TOPIC_ARN_ALERT: The topic to send exceptions.   \r\n\r\nKafka target environment variables:\r\nMSK_BOOTSTRAP_SRV: The URIs of the MSK cluster to publish messages. \r\n\r\nSNS target environment variables:\r\nSNS_TOPIC_ARN_EVENT: The topic to send docdb events.    \r\n\r\nS3 target environment variables:\r\nBUCKET_NAME: The name of the bucket that will save streamed data. \r\nBUCKET_PATH (optional): The path of the bucket that will save streamed data. \r\n\r\nElasticSearch target environment variables:\r\nELASTICSEARCH_URI: The URI of the Elasticsearch domain where data should be streamed.\r\n\r\nKinesis target environment variables:\r\nKINESIS_STREAM : The Kinesis Stream name to publish DocumentDB events.\r\n\r\nSQS target environment variables:\r\nSQS_QUERY_URL: The URL of the Amazon SQS queue to which a message is sent.\r\n\r\n\"\"\"<\/span>\r\n\r\n<span class=\"n\">db_client<\/span> <span class=\"o\">=<\/span> <span class=\"no\">None<\/span>                        <span class=\"c1\"># DocumentDB client - used as source <\/span>\r\n<span class=\"n\">s3_client<\/span> <span class=\"o\">=<\/span> <span class=\"no\">None<\/span>                        <span class=\"c1\"># S3 client - used as target        <\/span>\r\n<span class=\"n\">sns_client<\/span> <span class=\"o\">=<\/span> <span class=\"n\">boto3<\/span><span class=\"p\">.<\/span><span class=\"nf\">client<\/span><span class=\"p\">(<\/span><span class=\"s1\">'sns'<\/span><span class=\"p\">)<\/span>        <span class=\"c1\"># SNS client - for exception alerting purposes<\/span>\r\n<span class=\"n\">clientS3<\/span> <span class=\"o\">=<\/span> <span class=\"n\">boto3<\/span><span class=\"p\">.<\/span><span class=\"nf\">resource<\/span><span class=\"p\">(<\/span><span class=\"s1\">'s3'<\/span><span class=\"p\">)<\/span>         <span class=\"c1\"># S3 client - used to get the DocumentDB certificates<\/span>\r\n                                  \r\n<span class=\"n\">logger<\/span> <span class=\"o\">=<\/span> <span class=\"n\">logging<\/span><span class=\"p\">.<\/span><span class=\"nf\">getLogger<\/span><span class=\"p\">()<\/span>\r\n<span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">setLevel<\/span><span class=\"p\">(<\/span><span class=\"n\">logging<\/span><span class=\"o\">.<\/span><span class=\"no\">DEBUG<\/span><span class=\"p\">)<\/span>\r\n\r\n<span class=\"c1\"># The error code returned when data for the requested resume token has been deleted<\/span>\r\n<span class=\"no\">TOKEN_DATA_DELETED_CODE<\/span> <span class=\"o\">=<\/span> <span class=\"mi\">136<\/span>\r\n\r\n\r\n<span class=\"k\">def<\/span> <span class=\"nf\">get_credentials<\/span><span class=\"p\">():<\/span>\r\n    <span class=\"s2\">\"\"\"Retrieve credentials from the Secrets Manager service.\"\"\"<\/span>\r\n    <span class=\"n\">boto_session<\/span> <span class=\"o\">=<\/span> <span class=\"n\">boto3<\/span><span class=\"p\">.<\/span><span class=\"nf\">session<\/span><span class=\"o\">.<\/span><span class=\"no\">Session<\/span><span class=\"p\">()<\/span>\r\n\r\n    <span class=\"ss\">try:\r\n        <\/span><span class=\"n\">secret_name<\/span> <span class=\"o\">=<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'DOCUMENTDB_SECRET'<\/span><span class=\"p\">]<\/span>\r\n\r\n        <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Retrieving secret {} from Secrets Manger.'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">secret_name<\/span><span class=\"p\">))<\/span>\r\n\r\n        <span class=\"n\">secrets_client<\/span> <span class=\"o\">=<\/span> <span class=\"n\">boto_session<\/span><span class=\"p\">.<\/span><span class=\"nf\">client<\/span><span class=\"p\">(<\/span><span class=\"n\">service_name<\/span><span class=\"o\">=<\/span><span class=\"s1\">'secretsmanager'<\/span><span class=\"p\">,<\/span>\r\n                                             <span class=\"n\">region_name<\/span><span class=\"o\">=<\/span><span class=\"n\">boto_session<\/span><span class=\"p\">.<\/span><span class=\"nf\">region_name<\/span><span class=\"p\">)<\/span>\r\n        <span class=\"n\">secret_value<\/span> <span class=\"o\">=<\/span> <span class=\"n\">secrets_client<\/span><span class=\"p\">.<\/span><span class=\"nf\">get_secret_value<\/span><span class=\"p\">(<\/span><span class=\"no\">SecretId<\/span><span class=\"o\">=<\/span><span class=\"n\">secret_name<\/span><span class=\"p\">)<\/span>\r\n\r\n        <span class=\"n\">secret<\/span> <span class=\"o\">=<\/span> <span class=\"n\">secret_value<\/span><span class=\"p\">[<\/span><span class=\"s1\">'SecretString'<\/span><span class=\"p\">]<\/span>\r\n        <span class=\"n\">secret_json<\/span> <span class=\"o\">=<\/span> <span class=\"n\">json<\/span><span class=\"p\">.<\/span><span class=\"nf\">loads<\/span><span class=\"p\">(<\/span><span class=\"n\">secret<\/span><span class=\"p\">)<\/span>\r\n        <span class=\"n\">username<\/span> <span class=\"o\">=<\/span> <span class=\"n\">secret_json<\/span><span class=\"p\">[<\/span><span class=\"s1\">'username'<\/span><span class=\"p\">]<\/span>\r\n        <span class=\"n\">password<\/span> <span class=\"o\">=<\/span> <span class=\"n\">secret_json<\/span><span class=\"p\">[<\/span><span class=\"s1\">'password'<\/span><span class=\"p\">]<\/span>\r\n\r\n        <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Secret {} retrieved from Secrets Manger.'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">secret_name<\/span><span class=\"p\">))<\/span>\r\n\r\n        <span class=\"k\">return<\/span> <span class=\"p\">(<\/span><span class=\"n\">username<\/span><span class=\"p\">,<\/span> <span class=\"n\">password<\/span><span class=\"p\">)<\/span>\r\n\r\n    <span class=\"n\">except<\/span> <span class=\"no\">Exception<\/span> <span class=\"n\">as<\/span> <span class=\"ss\">ex:\r\n        <\/span><span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">error<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Failed to retrieve secret {}'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">secret_name<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"k\">raise<\/span>\r\n\r\n\r\n<span class=\"k\">def<\/span> <span class=\"nf\">get_db_client<\/span><span class=\"p\">():<\/span>\r\n    <span class=\"s2\">\"\"\"Return an authenticated connection to DocumentDB\"\"\"<\/span>\r\n    <span class=\"c1\"># Use a global variable so Lambda can reuse the persisted client on future invocations<\/span>\r\n    <span class=\"n\">global<\/span> <span class=\"n\">db_client<\/span>\r\n\r\n    <span class=\"k\">if<\/span> <span class=\"n\">db_client<\/span> <span class=\"n\">is<\/span> <span class=\"no\">None<\/span><span class=\"p\">:<\/span>\r\n        <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Creating new DocumentDB client.'<\/span><span class=\"p\">)<\/span>\r\n\r\n        <span class=\"ss\">try:\r\n            <\/span><span class=\"n\">cluster_uri<\/span> <span class=\"o\">=<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'DOCUMENTDB_URI'<\/span><span class=\"p\">]<\/span>\r\n            <span class=\"p\">(<\/span><span class=\"n\">username<\/span><span class=\"p\">,<\/span> <span class=\"n\">password<\/span><span class=\"p\">)<\/span> <span class=\"o\">=<\/span> <span class=\"n\">get_credentials<\/span><span class=\"p\">()<\/span>                   \r\n            <span class=\"n\">db_client<\/span> <span class=\"o\">=<\/span> <span class=\"no\">MongoClient<\/span><span class=\"p\">(<\/span><span class=\"n\">cluster_uri<\/span><span class=\"p\">,<\/span> <span class=\"n\">ssl<\/span><span class=\"o\">=<\/span><span class=\"no\">True<\/span><span class=\"p\">,<\/span> <span class=\"n\">retryWrites<\/span><span class=\"o\">=<\/span><span class=\"no\">False<\/span><span class=\"p\">,<\/span> <span class=\"n\">ssl_ca_certs<\/span><span class=\"o\">=<\/span><span class=\"s1\">'rds-combined-ca-bundle.pem'<\/span><span class=\"p\">)<\/span>\r\n            <span class=\"c1\"># force the client to connect<\/span>\r\n            <span class=\"n\">db_client<\/span><span class=\"p\">.<\/span><span class=\"nf\">admin<\/span><span class=\"p\">.<\/span><span class=\"nf\">command<\/span><span class=\"p\">(<\/span><span class=\"s1\">'ismaster'<\/span><span class=\"p\">)<\/span>\r\n            <span class=\"n\">db_client<\/span><span class=\"p\">[<\/span><span class=\"s2\">\"admin\"<\/span><span class=\"p\">].<\/span><span class=\"nf\">authenticate<\/span><span class=\"p\">(<\/span><span class=\"nb\">name<\/span><span class=\"o\">=<\/span><span class=\"n\">username<\/span><span class=\"p\">,<\/span> <span class=\"n\">password<\/span><span class=\"o\">=<\/span><span class=\"n\">password<\/span><span class=\"p\">)<\/span>\r\n\r\n            <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Successfully created new DocumentDB client.'<\/span><span class=\"p\">)<\/span>\r\n        <span class=\"n\">except<\/span> <span class=\"no\">Exception<\/span> <span class=\"n\">as<\/span> <span class=\"ss\">ex:\r\n            <\/span><span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">error<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Failed to create new DocumentDB client: {}'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n            <span class=\"n\">send_sns_alert<\/span><span class=\"p\">(<\/span><span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n            <span class=\"k\">raise<\/span>\r\n\r\n    <span class=\"k\">return<\/span> <span class=\"n\">db_client<\/span>\r\n\r\n\r\n<span class=\"k\">def<\/span> <span class=\"nf\">get_state_collection_client<\/span><span class=\"p\">():<\/span>\r\n    <span class=\"s2\">\"\"\"Return a DocumentDB client for the collection in which we store processing state.\"\"\"<\/span>\r\n\r\n    <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Creating state_collection_client.'<\/span><span class=\"p\">)<\/span>\r\n    <span class=\"ss\">try:\r\n        <\/span><span class=\"n\">db_client<\/span> <span class=\"o\">=<\/span> <span class=\"n\">get_db_client<\/span><span class=\"p\">()<\/span>\r\n        <span class=\"n\">state_db_name<\/span> <span class=\"o\">=<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'STATE_DB'<\/span><span class=\"p\">]<\/span>\r\n        <span class=\"n\">state_collection_name<\/span> <span class=\"o\">=<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'STATE_COLLECTION'<\/span><span class=\"p\">]<\/span>\r\n        <span class=\"n\">state_collection<\/span> <span class=\"o\">=<\/span> <span class=\"n\">db_client<\/span><span class=\"p\">[<\/span><span class=\"n\">state_db_name<\/span><span class=\"p\">][<\/span><span class=\"n\">state_collection_name<\/span><span class=\"p\">]<\/span>\r\n    <span class=\"n\">except<\/span> <span class=\"no\">Exception<\/span> <span class=\"n\">as<\/span> <span class=\"ss\">ex:\r\n        <\/span><span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">error<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Failed to create new state collection client: {}'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"n\">send_sns_alert<\/span><span class=\"p\">(<\/span><span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"k\">raise<\/span>\r\n\r\n    <span class=\"k\">return<\/span> <span class=\"n\">state_collection<\/span>\r\n\r\n\r\n<span class=\"k\">def<\/span> <span class=\"nf\">get_last_processed_id<\/span><span class=\"p\">():<\/span>\r\n    <span class=\"s2\">\"\"\"Return the resume token corresponding to the last successfully processed change event.\"\"\"<\/span>\r\n    <span class=\"n\">last_processed_id<\/span> <span class=\"o\">=<\/span> <span class=\"no\">None<\/span>\r\n    <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Returning last processed id.'<\/span><span class=\"p\">)<\/span>\r\n    <span class=\"ss\">try:\r\n        <\/span><span class=\"n\">state_collection<\/span> <span class=\"o\">=<\/span> <span class=\"n\">get_state_collection_client<\/span><span class=\"p\">()<\/span>\r\n        <span class=\"k\">if<\/span> <span class=\"s2\">\"WATCHED_COLLECTION_NAME\"<\/span> <span class=\"k\">in<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">:<\/span>\r\n            <span class=\"n\">state_doc<\/span> <span class=\"o\">=<\/span> <span class=\"n\">state_collection<\/span><span class=\"p\">.<\/span><span class=\"nf\">find_one<\/span><span class=\"p\">({<\/span><span class=\"s1\">'currentState'<\/span><span class=\"p\">:<\/span> <span class=\"no\">True<\/span><span class=\"p\">,<\/span> <span class=\"s1\">'dbWatched'<\/span><span class=\"p\">:<\/span> <span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'WATCHED_DB_NAME'<\/span><span class=\"p\">]),<\/span> \r\n                <span class=\"s1\">'collectionWatched'<\/span><span class=\"p\">:<\/span> <span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'WATCHED_COLLECTION_NAME'<\/span><span class=\"p\">]),<\/span> <span class=\"s1\">'db_level'<\/span><span class=\"p\">:<\/span> <span class=\"no\">False<\/span><span class=\"p\">})<\/span>\r\n        <span class=\"ss\">else:\r\n            <\/span><span class=\"n\">state_doc<\/span> <span class=\"o\">=<\/span> <span class=\"n\">state_collection<\/span><span class=\"p\">.<\/span><span class=\"nf\">find_one<\/span><span class=\"p\">({<\/span><span class=\"s1\">'currentState'<\/span><span class=\"p\">:<\/span> <span class=\"no\">True<\/span><span class=\"p\">,<\/span> <span class=\"s1\">'db_level'<\/span><span class=\"p\">:<\/span> <span class=\"no\">True<\/span><span class=\"p\">,<\/span> \r\n                <span class=\"s1\">'dbWatched'<\/span><span class=\"p\">:<\/span> <span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'WATCHED_DB_NAME'<\/span><span class=\"p\">])})<\/span>\r\n           \r\n        <span class=\"k\">if<\/span> <span class=\"n\">state_doc<\/span> <span class=\"n\">is<\/span> <span class=\"ow\">not<\/span> <span class=\"no\">None<\/span><span class=\"p\">:<\/span>\r\n            <span class=\"k\">if<\/span> <span class=\"s1\">'lastProcessed'<\/span> <span class=\"k\">in<\/span> <span class=\"ss\">state_doc: \r\n                <\/span><span class=\"n\">last_processed_id<\/span> <span class=\"o\">=<\/span> <span class=\"n\">state_doc<\/span><span class=\"p\">[<\/span><span class=\"s1\">'lastProcessed'<\/span><span class=\"p\">]<\/span>\r\n        <span class=\"ss\">else:\r\n            <\/span><span class=\"k\">if<\/span> <span class=\"s2\">\"WATCHED_COLLECTION_NAME\"<\/span> <span class=\"k\">in<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">:<\/span>\r\n                <span class=\"n\">state_collection<\/span><span class=\"p\">.<\/span><span class=\"nf\">insert_one<\/span><span class=\"p\">({<\/span><span class=\"s1\">'dbWatched'<\/span><span class=\"p\">:<\/span> <span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'WATCHED_DB_NAME'<\/span><span class=\"p\">]),<\/span>\r\n                    <span class=\"s1\">'collectionWatched'<\/span><span class=\"p\">:<\/span> <span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'WATCHED_COLLECTION_NAME'<\/span><span class=\"p\">]),<\/span> <span class=\"s1\">'currentState'<\/span><span class=\"p\">:<\/span> <span class=\"no\">True<\/span><span class=\"p\">,<\/span> <span class=\"s1\">'db_level'<\/span><span class=\"p\">:<\/span> <span class=\"no\">False<\/span><span class=\"p\">})<\/span>\r\n            <span class=\"ss\">else:\r\n                <\/span><span class=\"n\">state_collection<\/span><span class=\"p\">.<\/span><span class=\"nf\">insert_one<\/span><span class=\"p\">({<\/span><span class=\"s1\">'dbWatched'<\/span><span class=\"p\">:<\/span> <span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'WATCHED_DB_NAME'<\/span><span class=\"p\">]),<\/span> <span class=\"s1\">'currentState'<\/span><span class=\"p\">:<\/span> <span class=\"no\">True<\/span><span class=\"p\">,<\/span> \r\n                    <span class=\"s1\">'db_level'<\/span><span class=\"p\">:<\/span> <span class=\"no\">True<\/span><span class=\"p\">})<\/span>\r\n\r\n    <span class=\"n\">except<\/span> <span class=\"no\">Exception<\/span> <span class=\"n\">as<\/span> <span class=\"ss\">ex:\r\n        <\/span><span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">error<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Failed to return last processed id: {}'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"n\">send_sns_alert<\/span><span class=\"p\">(<\/span><span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"k\">raise<\/span>\r\n\r\n    <span class=\"k\">return<\/span> <span class=\"n\">last_processed_id<\/span>\r\n\r\n\r\n<span class=\"k\">def<\/span> <span class=\"nf\">store_last_processed_id<\/span><span class=\"p\">(<\/span><span class=\"n\">resume_token<\/span><span class=\"p\">):<\/span>\r\n    <span class=\"s2\">\"\"\"Store the resume token corresponding to the last successfully processed change event.\"\"\"<\/span>\r\n\r\n    <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Storing last processed id.'<\/span><span class=\"p\">)<\/span>\r\n    <span class=\"ss\">try:\r\n        <\/span><span class=\"n\">state_collection<\/span> <span class=\"o\">=<\/span> <span class=\"n\">get_state_collection_client<\/span><span class=\"p\">()<\/span>\r\n        <span class=\"k\">if<\/span> <span class=\"s2\">\"WATCHED_COLLECTION_NAME\"<\/span> <span class=\"k\">in<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">:<\/span>\r\n            <span class=\"n\">state_collection<\/span><span class=\"p\">.<\/span><span class=\"nf\">update_one<\/span><span class=\"p\">({<\/span><span class=\"s1\">'dbWatched'<\/span><span class=\"p\">:<\/span> <span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'WATCHED_DB_NAME'<\/span><span class=\"p\">]),<\/span> \r\n                <span class=\"s1\">'collectionWatched'<\/span><span class=\"p\">:<\/span> <span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'WATCHED_COLLECTION_NAME'<\/span><span class=\"p\">])},{<\/span><span class=\"s1\">'$set'<\/span><span class=\"p\">:<\/span> <span class=\"p\">{<\/span><span class=\"s1\">'lastProcessed'<\/span><span class=\"p\">:<\/span> <span class=\"n\">resume_token<\/span><span class=\"p\">}})<\/span>\r\n        <span class=\"ss\">else:\r\n            <\/span><span class=\"n\">state_collection<\/span><span class=\"p\">.<\/span><span class=\"nf\">update_one<\/span><span class=\"p\">({<\/span><span class=\"s1\">'dbWatched'<\/span><span class=\"p\">:<\/span> <span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'WATCHED_DB_NAME'<\/span><span class=\"p\">]),<\/span> <span class=\"s1\">'db_level'<\/span><span class=\"p\">:<\/span> <span class=\"no\">True<\/span><span class=\"p\">,<\/span> <span class=\"p\">},<\/span>\r\n                <span class=\"p\">{<\/span><span class=\"s1\">'$set'<\/span><span class=\"p\">:<\/span> <span class=\"p\">{<\/span><span class=\"s1\">'lastProcessed'<\/span><span class=\"p\">:<\/span> <span class=\"n\">resume_token<\/span><span class=\"p\">}})<\/span>\r\n\r\n    <span class=\"n\">except<\/span> <span class=\"no\">Exception<\/span> <span class=\"n\">as<\/span> <span class=\"ss\">ex:\r\n        <\/span><span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">error<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Failed to store last processed id: {}'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"n\">send_sns_alert<\/span><span class=\"p\">(<\/span><span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"k\">raise<\/span>\r\n\r\n\r\n<span class=\"k\">def<\/span> <span class=\"nf\">send_sns_alert<\/span><span class=\"p\">(<\/span><span class=\"n\">message<\/span><span class=\"p\">):<\/span>\r\n    <span class=\"s2\">\"\"\"send an SNS alert\"\"\"<\/span>\r\n    <span class=\"ss\">try:\r\n        <\/span><span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Sending SNS alert.'<\/span><span class=\"p\">)<\/span>\r\n        <span class=\"n\">response<\/span> <span class=\"o\">=<\/span> <span class=\"n\">sns_client<\/span><span class=\"p\">.<\/span><span class=\"nf\">publish<\/span><span class=\"p\">(<\/span>\r\n            <span class=\"no\">TopicArn<\/span><span class=\"o\">=<\/span><span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'SNS_TOPIC_ARN_ALERT'<\/span><span class=\"p\">],<\/span>\r\n            <span class=\"no\">Message<\/span><span class=\"o\">=<\/span><span class=\"n\">message<\/span><span class=\"p\">,<\/span>\r\n            <span class=\"no\">Subject<\/span><span class=\"o\">=<\/span><span class=\"s1\">'Document DB Replication Alarm'<\/span><span class=\"p\">,<\/span>\r\n            <span class=\"no\">MessageStructure<\/span><span class=\"o\">=<\/span><span class=\"s1\">'default'<\/span>\r\n        <span class=\"p\">)<\/span>\r\n    <span class=\"n\">except<\/span> <span class=\"no\">Exception<\/span> <span class=\"n\">as<\/span> <span class=\"ss\">ex:\r\n        <\/span><span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">error<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Exception in publishing alert to SNS: {}'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"n\">send_sns_alert<\/span><span class=\"p\">(<\/span><span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"k\">raise<\/span>\r\n\r\n\r\n<span class=\"k\">def<\/span> <span class=\"nf\">getDocDbCertificate<\/span><span class=\"p\">():<\/span>\r\n    <span class=\"s2\">\"\"\"download the current docdb certificate\"\"\"<\/span>\r\n    <span class=\"ss\">try:\r\n        <\/span><span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Getting DocumentDB certificate from S3.'<\/span><span class=\"p\">)<\/span>\r\n        <span class=\"n\">clientS3<\/span><span class=\"o\">.<\/span><span class=\"no\">Bucket<\/span><span class=\"p\">(<\/span><span class=\"s1\">'rds-downloads'<\/span><span class=\"p\">).<\/span><span class=\"nf\">download_file<\/span><span class=\"p\">(<\/span><span class=\"s1\">'rds-combined-ca-bundle.pem'<\/span><span class=\"p\">,<\/span> <span class=\"s1\">'\/tmp\/rds-combined-ca-bundle.pem'<\/span><span class=\"p\">)<\/span>\r\n    <span class=\"n\">except<\/span> <span class=\"no\">Exception<\/span> <span class=\"n\">as<\/span> <span class=\"ss\">ex:\r\n        <\/span><span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">error<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Exception in publishing message to Kinesis: {}'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"n\">send_sns_alert<\/span><span class=\"p\">(<\/span><span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"k\">raise<\/span>\r\n\r\n\r\n<span class=\"k\">def<\/span> <span class=\"nf\">insertCanary<\/span><span class=\"p\">():<\/span>\r\n    <span class=\"s2\">\"\"\"Inserts a canary event for change stream activation\"\"\"<\/span>\r\n    \r\n    <span class=\"n\">canary_record<\/span> <span class=\"o\">=<\/span> <span class=\"no\">None<\/span>\r\n\r\n    <span class=\"ss\">try:\r\n        <\/span><span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Inserting canary.'<\/span><span class=\"p\">)<\/span>\r\n        <span class=\"n\">db_client<\/span> <span class=\"o\">=<\/span> <span class=\"n\">get_db_client<\/span><span class=\"p\">()<\/span>\r\n        <span class=\"n\">watched_db<\/span> <span class=\"o\">=<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'WATCHED_DB_NAME'<\/span><span class=\"p\">]<\/span>\r\n\r\n        <span class=\"k\">if<\/span> <span class=\"s2\">\"WATCHED_COLLECTION_NAME\"<\/span> <span class=\"k\">in<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">:<\/span>\r\n            <span class=\"n\">watched_collection<\/span> <span class=\"o\">=<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'WATCHED_COLLECTION_NAME'<\/span><span class=\"p\">]<\/span>\r\n        <span class=\"ss\">else:\r\n            <\/span><span class=\"n\">watched_collection<\/span> <span class=\"o\">=<\/span> <span class=\"s1\">'canary-collection'<\/span>\r\n\r\n        <span class=\"n\">collection_client<\/span> <span class=\"o\">=<\/span> <span class=\"n\">db_client<\/span><span class=\"p\">[<\/span><span class=\"n\">watched_db<\/span><span class=\"p\">][<\/span><span class=\"n\">watched_collection<\/span><span class=\"p\">]<\/span>\r\n\r\n        <span class=\"n\">canary_record<\/span> <span class=\"o\">=<\/span> <span class=\"n\">collection_client<\/span><span class=\"p\">.<\/span><span class=\"nf\">insert_one<\/span><span class=\"p\">({<\/span> <span class=\"s2\">\"op_canary\"<\/span><span class=\"p\">:<\/span> <span class=\"s2\">\"canary\"<\/span> <span class=\"p\">})<\/span>\r\n        <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Canary inserted.'<\/span><span class=\"p\">)<\/span>\r\n\r\n    <span class=\"n\">except<\/span> <span class=\"no\">Exception<\/span> <span class=\"n\">as<\/span> <span class=\"ss\">ex:\r\n        <\/span><span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">error<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Exception in inserting canary: {}'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"n\">send_sns_alert<\/span><span class=\"p\">(<\/span><span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"k\">raise<\/span>\r\n\r\n    <span class=\"k\">return<\/span> <span class=\"n\">canary_record<\/span>\r\n\r\n\r\n<span class=\"k\">def<\/span> <span class=\"nf\">deleteCanary<\/span><span class=\"p\">():<\/span>\r\n    <span class=\"s2\">\"\"\"Deletes a canary event for change stream activation\"\"\"<\/span>\r\n    \r\n    <span class=\"ss\">try:\r\n        <\/span><span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Deleting canary.'<\/span><span class=\"p\">)<\/span>\r\n        <span class=\"n\">db_client<\/span> <span class=\"o\">=<\/span> <span class=\"n\">get_db_client<\/span><span class=\"p\">()<\/span>\r\n        <span class=\"n\">watched_db<\/span> <span class=\"o\">=<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'WATCHED_DB_NAME'<\/span><span class=\"p\">]<\/span>\r\n\r\n        <span class=\"k\">if<\/span> <span class=\"s2\">\"WATCHED_COLLECTION_NAME\"<\/span> <span class=\"k\">in<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">:<\/span>\r\n            <span class=\"n\">watched_collection<\/span> <span class=\"o\">=<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'WATCHED_COLLECTION_NAME'<\/span><span class=\"p\">]<\/span>\r\n        <span class=\"ss\">else:\r\n            <\/span><span class=\"n\">watched_collection<\/span> <span class=\"o\">=<\/span> <span class=\"s1\">'canary-collection'<\/span>\r\n\r\n        <span class=\"n\">collection_client<\/span> <span class=\"o\">=<\/span> <span class=\"n\">db_client<\/span><span class=\"p\">[<\/span><span class=\"n\">watched_db<\/span><span class=\"p\">][<\/span><span class=\"n\">watched_collection<\/span><span class=\"p\">]<\/span>\r\n        <span class=\"n\">collection_client<\/span><span class=\"p\">.<\/span><span class=\"nf\">delete_one<\/span><span class=\"p\">({<\/span> <span class=\"s2\">\"op_canary\"<\/span><span class=\"p\">:<\/span> <span class=\"s2\">\"canary\"<\/span> <span class=\"p\">})<\/span>\r\n        <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Canary deleted.'<\/span><span class=\"p\">)<\/span>\r\n    \r\n    <span class=\"n\">except<\/span> <span class=\"no\">Exception<\/span> <span class=\"n\">as<\/span> <span class=\"ss\">ex:\r\n        <\/span><span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">error<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Exception in deleting canary: {}'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"n\">send_sns_alert<\/span><span class=\"p\">(<\/span><span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"k\">raise<\/span>\r\n\r\n\r\n<span class=\"k\">def<\/span> <span class=\"nf\">put_s3_event<\/span><span class=\"p\">(<\/span><span class=\"n\">event<\/span><span class=\"p\">,<\/span> <span class=\"n\">database<\/span><span class=\"p\">,<\/span> <span class=\"n\">collection<\/span><span class=\"p\">,<\/span> <span class=\"n\">doc_id<\/span><span class=\"p\">):<\/span>\r\n    <span class=\"s2\">\"\"\"send event to S3\"\"\"<\/span>\r\n    <span class=\"c1\"># Use a global variable so Lambda can reuse the persisted client on future invocations<\/span>\r\n    <span class=\"n\">global<\/span> <span class=\"n\">s3_client<\/span>\r\n\r\n    <span class=\"k\">if<\/span> <span class=\"n\">s3_client<\/span> <span class=\"n\">is<\/span> <span class=\"no\">None<\/span><span class=\"p\">:<\/span>\r\n        <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Creating new S3 client.'<\/span><span class=\"p\">)<\/span>\r\n        <span class=\"n\">s3_client<\/span> <span class=\"o\">=<\/span> <span class=\"n\">boto3<\/span><span class=\"p\">.<\/span><span class=\"nf\">resource<\/span><span class=\"p\">(<\/span><span class=\"s1\">'s3'<\/span><span class=\"p\">)<\/span>  \r\n\r\n    <span class=\"ss\">try:\r\n        <\/span><span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Publishing message to S3.'<\/span><span class=\"p\">)<\/span> <span class=\"c1\">#, str(os.environ['BUCKET_PATH'])<\/span>\r\n        <span class=\"k\">if<\/span> <span class=\"s2\">\"BUCKET_PATH\"<\/span> <span class=\"k\">in<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">:<\/span>\r\n            <span class=\"n\">s3_client<\/span><span class=\"o\">.<\/span><span class=\"no\">Object<\/span><span class=\"p\">(<\/span><span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'BUCKET_NAME'<\/span><span class=\"p\">],<\/span> <span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'BUCKET_PATH'<\/span><span class=\"p\">])<\/span> <span class=\"o\">+<\/span> <span class=\"s1\">'\/'<\/span> <span class=\"o\">+<\/span> <span class=\"n\">database<\/span> <span class=\"o\">+<\/span> <span class=\"s1\">'\/'<\/span> <span class=\"o\">+<\/span>\r\n                <span class=\"n\">collection<\/span> <span class=\"o\">+<\/span> <span class=\"s1\">'\/'<\/span> <span class=\"o\">+<\/span> <span class=\"n\">datetime<\/span><span class=\"p\">.<\/span><span class=\"nf\">datetime<\/span><span class=\"p\">.<\/span><span class=\"nf\">now<\/span><span class=\"p\">().<\/span><span class=\"nf\">strftime<\/span><span class=\"p\">(<\/span><span class=\"s1\">'%Y%m%d'<\/span><span class=\"p\">)<\/span> <span class=\"o\">+<\/span> <span class=\"s1\">'\/'<\/span> <span class=\"o\">+<\/span> <span class=\"n\">doc_id<\/span><span class=\"p\">).<\/span><span class=\"nf\">put<\/span><span class=\"p\">(<\/span><span class=\"no\">Body<\/span><span class=\"o\">=<\/span><span class=\"n\">event<\/span><span class=\"p\">)<\/span>\r\n        <span class=\"ss\">else: \r\n            <\/span><span class=\"n\">s3_client<\/span><span class=\"o\">.<\/span><span class=\"no\">Object<\/span><span class=\"p\">(<\/span><span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'BUCKET_NAME'<\/span><span class=\"p\">],<\/span> <span class=\"n\">database<\/span> <span class=\"o\">+<\/span> <span class=\"s1\">'\/'<\/span> <span class=\"o\">+<\/span> <span class=\"n\">collection<\/span> <span class=\"o\">+<\/span> <span class=\"s1\">'\/'<\/span> <span class=\"o\">+<\/span> \r\n                <span class=\"n\">datetime<\/span><span class=\"p\">.<\/span><span class=\"nf\">datetime<\/span><span class=\"p\">.<\/span><span class=\"nf\">now<\/span><span class=\"p\">().<\/span><span class=\"nf\">strftime<\/span><span class=\"p\">(<\/span><span class=\"s1\">'%Y%m%d'<\/span><span class=\"p\">)<\/span> <span class=\"o\">+<\/span> <span class=\"s1\">'\/'<\/span> <span class=\"o\">+<\/span> <span class=\"n\">doc_id<\/span><span class=\"p\">).<\/span><span class=\"nf\">put<\/span><span class=\"p\">(<\/span><span class=\"no\">Body<\/span><span class=\"o\">=<\/span><span class=\"n\">event<\/span><span class=\"p\">)<\/span>\r\n\r\n    <span class=\"n\">except<\/span> <span class=\"no\">Exception<\/span> <span class=\"n\">as<\/span> <span class=\"ss\">ex:\r\n        <\/span><span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">error<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Exception in publishing message to S3: {}'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"n\">send_sns_alert<\/span><span class=\"p\">(<\/span><span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"k\">raise<\/span>\r\n\r\n<span class=\"k\">def<\/span> <span class=\"nf\">lambda_handler<\/span><span class=\"p\">(<\/span><span class=\"n\">event<\/span><span class=\"p\">,<\/span> <span class=\"n\">context<\/span><span class=\"p\">):<\/span>\r\n    <span class=\"s2\">\"\"\"Read any new events from DocumentDB and apply them to an streaming\/datastore endpoint.\"\"\"<\/span>\r\n    \r\n    <span class=\"n\">events_processed<\/span> <span class=\"o\">=<\/span> <span class=\"mi\">0<\/span>\r\n    <span class=\"n\">canary_record<\/span> <span class=\"o\">=<\/span> <span class=\"no\">None<\/span>\r\n    <span class=\"n\">watcher<\/span> <span class=\"o\">=<\/span> <span class=\"no\">None<\/span>\r\n    <span class=\"n\">folder<\/span> <span class=\"o\">=<\/span> <span class=\"no\">None<\/span>\r\n    <span class=\"n\">filename<\/span> <span class=\"o\">=<\/span> <span class=\"no\">None<\/span>\r\n    <span class=\"n\">kafka_client<\/span> <span class=\"o\">=<\/span> <span class=\"no\">None<\/span>\r\n    <span class=\"c1\"># getDocDbCertificate()<\/span>\r\n\r\n    <span class=\"ss\">try:\r\n        \r\n        <\/span><span class=\"c1\"># DocumentDB watched collection set up<\/span>\r\n        <span class=\"n\">db_client<\/span> <span class=\"o\">=<\/span> <span class=\"n\">get_db_client<\/span><span class=\"p\">()<\/span>\r\n        <span class=\"n\">watched_db<\/span> <span class=\"o\">=<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'WATCHED_DB_NAME'<\/span><span class=\"p\">]<\/span>\r\n        <span class=\"k\">if<\/span> <span class=\"s2\">\"WATCHED_COLLECTION_NAME\"<\/span> <span class=\"k\">in<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">:<\/span>\r\n            <span class=\"n\">watched_collection<\/span> <span class=\"o\">=<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'WATCHED_COLLECTION_NAME'<\/span><span class=\"p\">]<\/span>\r\n            <span class=\"n\">watcher<\/span> <span class=\"o\">=<\/span> <span class=\"n\">db_client<\/span><span class=\"p\">[<\/span><span class=\"n\">watched_db<\/span><span class=\"p\">][<\/span><span class=\"n\">watched_collection<\/span><span class=\"p\">]<\/span>\r\n        <span class=\"ss\">else: \r\n            <\/span><span class=\"n\">watcher<\/span> <span class=\"o\">=<\/span> <span class=\"n\">db_client<\/span><span class=\"p\">[<\/span><span class=\"n\">watched_db<\/span><span class=\"p\">]<\/span>\r\n        <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Watching collection {}'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">watcher<\/span><span class=\"p\">))<\/span>\r\n\r\n        <span class=\"c1\"># DocumentDB sync set up<\/span>\r\n        <span class=\"n\">state_sync_count<\/span> <span class=\"o\">=<\/span> <span class=\"n\">int<\/span><span class=\"p\">(<\/span><span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'Iterations_per_sync'<\/span><span class=\"p\">])<\/span>\r\n        <span class=\"n\">last_processed_id<\/span> <span class=\"o\">=<\/span> <span class=\"n\">get_last_processed_id<\/span><span class=\"p\">()<\/span>\r\n        <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s2\">\"last_processed_id: {}\"<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">last_processed_id<\/span><span class=\"p\">))<\/span>\r\n\r\n        <span class=\"n\">with<\/span> <span class=\"n\">watcher<\/span><span class=\"p\">.<\/span><span class=\"nf\">watch<\/span><span class=\"p\">(<\/span><span class=\"n\">full_document<\/span><span class=\"o\">=<\/span><span class=\"s1\">'updateLookup'<\/span><span class=\"p\">,<\/span> <span class=\"n\">resume_after<\/span><span class=\"o\">=<\/span><span class=\"n\">last_processed_id<\/span><span class=\"p\">)<\/span> <span class=\"n\">as<\/span> <span class=\"ss\">change_stream:\r\n            <\/span><span class=\"n\">i<\/span> <span class=\"o\">=<\/span> <span class=\"mi\">0<\/span>\r\n\r\n            <span class=\"k\">if<\/span> <span class=\"n\">last_processed_id<\/span> <span class=\"n\">is<\/span> <span class=\"no\">None<\/span><span class=\"p\">:<\/span>\r\n                <span class=\"n\">canary_record<\/span> <span class=\"o\">=<\/span> <span class=\"n\">insertCanary<\/span><span class=\"p\">()<\/span>\r\n                <span class=\"n\">deleteCanary<\/span><span class=\"p\">()<\/span>\r\n\r\n            <span class=\"k\">while<\/span> <span class=\"n\">change_stream<\/span><span class=\"p\">.<\/span><span class=\"nf\">alive<\/span> <span class=\"ow\">and<\/span> <span class=\"n\">i<\/span> <span class=\"o\">&lt;<\/span> <span class=\"n\">int<\/span><span class=\"p\">(<\/span><span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">[<\/span><span class=\"s1\">'Documents_per_run'<\/span><span class=\"p\">]):<\/span>\r\n            \r\n                <span class=\"n\">i<\/span> <span class=\"o\">+=<\/span> <span class=\"mi\">1<\/span>\r\n                <span class=\"n\">change_event<\/span> <span class=\"o\">=<\/span> <span class=\"n\">change_stream<\/span><span class=\"p\">.<\/span><span class=\"nf\">try_next<\/span><span class=\"p\">()<\/span>\r\n                <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Event: {}'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">change_event<\/span><span class=\"p\">))<\/span>\r\n                \r\n                <span class=\"k\">if<\/span> <span class=\"n\">last_processed_id<\/span> <span class=\"n\">is<\/span> <span class=\"no\">None<\/span><span class=\"p\">:<\/span>\r\n                    <span class=\"k\">if<\/span> <span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'operationType'<\/span><span class=\"p\">]<\/span> <span class=\"o\">==<\/span> <span class=\"s1\">'delete'<\/span><span class=\"p\">:<\/span>\r\n                        <span class=\"n\">store_last_processed_id<\/span><span class=\"p\">(<\/span><span class=\"n\">change_stream<\/span><span class=\"p\">.<\/span><span class=\"nf\">resume_token<\/span><span class=\"p\">)<\/span>\r\n                        <span class=\"n\">last_processed_id<\/span> <span class=\"o\">=<\/span> <span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'_id'<\/span><span class=\"p\">][<\/span><span class=\"s1\">'_data'<\/span><span class=\"p\">]<\/span>\r\n                    <span class=\"n\">continue<\/span>\r\n                \r\n                <span class=\"k\">if<\/span> <span class=\"n\">change_event<\/span> <span class=\"n\">is<\/span> <span class=\"no\">None<\/span><span class=\"p\">:<\/span>\r\n                        <span class=\"k\">break<\/span>\r\n                <span class=\"ss\">else:\r\n                    <\/span><span class=\"n\">op_type<\/span> <span class=\"o\">=<\/span> <span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'operationType'<\/span><span class=\"p\">]<\/span>\r\n                    <span class=\"n\">op_id<\/span> <span class=\"o\">=<\/span> <span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'_id'<\/span><span class=\"p\">][<\/span><span class=\"s1\">'_data'<\/span><span class=\"p\">]<\/span>\r\n\r\n                    <span class=\"k\">if<\/span> <span class=\"n\">op_type<\/span> <span class=\"k\">in<\/span> <span class=\"p\">[<\/span><span class=\"s1\">'insert'<\/span><span class=\"p\">,<\/span> <span class=\"s1\">'update'<\/span><span class=\"p\">]:<\/span>             \r\n                        <span class=\"n\">doc_body<\/span> <span class=\"o\">=<\/span> <span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'fullDocument'<\/span><span class=\"p\">]<\/span>\r\n                        <span class=\"n\">doc_id<\/span> <span class=\"o\">=<\/span> <span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">doc_body<\/span><span class=\"p\">.<\/span><span class=\"nf\">pop<\/span><span class=\"p\">(<\/span><span class=\"s2\">\"_id\"<\/span><span class=\"p\">,<\/span> <span class=\"no\">None<\/span><span class=\"p\">))<\/span>\r\n                        <span class=\"n\">readable<\/span> <span class=\"o\">=<\/span> <span class=\"n\">datetime<\/span><span class=\"p\">.<\/span><span class=\"nf\">datetime<\/span><span class=\"p\">.<\/span><span class=\"nf\">fromtimestamp<\/span><span class=\"p\">(<\/span><span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'clusterTime'<\/span><span class=\"p\">].<\/span><span class=\"nf\">time<\/span><span class=\"p\">).<\/span><span class=\"nf\">isoformat<\/span><span class=\"p\">()<\/span>\r\n                        <span class=\"c1\">######## Uncomment the following line if you want to add operation metadata fields to the document event. <\/span>\r\n                        <span class=\"n\">doc_body<\/span><span class=\"p\">.<\/span><span class=\"nf\">update<\/span><span class=\"p\">({<\/span><span class=\"s1\">'action'<\/span><span class=\"ss\">:op_type<\/span><span class=\"p\">,<\/span><span class=\"s1\">'timestamp'<\/span><span class=\"ss\">:str<\/span><span class=\"p\">(<\/span><span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'clusterTime'<\/span><span class=\"p\">].<\/span><span class=\"nf\">time<\/span><span class=\"p\">),<\/span><span class=\"s1\">'timestampReadable'<\/span><span class=\"ss\">:str<\/span><span class=\"p\">(<\/span><span class=\"n\">readable<\/span><span class=\"p\">)})<\/span>\r\n                        <span class=\"c1\">######## Uncomment the following line if you want to add db and coll metadata fields to the document event. <\/span>\r\n                        <span class=\"n\">doc_body<\/span><span class=\"p\">.<\/span><span class=\"nf\">update<\/span><span class=\"p\">({<\/span><span class=\"s1\">'db'<\/span><span class=\"ss\">:str<\/span><span class=\"p\">(<\/span><span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'ns'<\/span><span class=\"p\">][<\/span><span class=\"s1\">'db'<\/span><span class=\"p\">]),<\/span><span class=\"s1\">'coll'<\/span><span class=\"ss\">:str<\/span><span class=\"p\">(<\/span><span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'ns'<\/span><span class=\"p\">][<\/span><span class=\"s1\">'coll'<\/span><span class=\"p\">])})<\/span>\r\n                        <span class=\"n\">payload<\/span> <span class=\"o\">=<\/span> <span class=\"p\">{<\/span><span class=\"s1\">'_id'<\/span><span class=\"ss\">:doc_id<\/span><span class=\"p\">}<\/span>\r\n                        <span class=\"n\">payload<\/span><span class=\"p\">.<\/span><span class=\"nf\">update<\/span><span class=\"p\">(<\/span><span class=\"n\">doc_body<\/span><span class=\"p\">)<\/span>\r\n\r\n                        <span class=\"c1\"># Append event for S3 <\/span>\r\n                        <span class=\"k\">if<\/span> <span class=\"s2\">\"BUCKET_NAME\"<\/span> <span class=\"k\">in<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">:<\/span>\r\n                            <span class=\"n\">put_s3_event<\/span><span class=\"p\">(<\/span><span class=\"n\">json_util<\/span><span class=\"p\">.<\/span><span class=\"nf\">dumps<\/span><span class=\"p\">(<\/span><span class=\"n\">payload<\/span><span class=\"p\">),<\/span> <span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'ns'<\/span><span class=\"p\">][<\/span><span class=\"s1\">'db'<\/span><span class=\"p\">]),<\/span> <span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'ns'<\/span><span class=\"p\">][<\/span><span class=\"s1\">'coll'<\/span><span class=\"p\">]),<\/span><span class=\"n\">op_id<\/span><span class=\"p\">)<\/span>\r\n                        \r\n                        <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Processed event ID {}'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">op_id<\/span><span class=\"p\">))<\/span>\r\n\r\n                    <span class=\"k\">if<\/span> <span class=\"n\">op_type<\/span> <span class=\"o\">==<\/span> <span class=\"s1\">'delete'<\/span><span class=\"p\">:<\/span>\r\n                        <span class=\"n\">doc_id<\/span> <span class=\"o\">=<\/span> <span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'documentKey'<\/span><span class=\"p\">][<\/span><span class=\"s1\">'_id'<\/span><span class=\"p\">])<\/span>\r\n                        <span class=\"n\">readable<\/span> <span class=\"o\">=<\/span> <span class=\"n\">datetime<\/span><span class=\"p\">.<\/span><span class=\"nf\">datetime<\/span><span class=\"p\">.<\/span><span class=\"nf\">fromtimestamp<\/span><span class=\"p\">(<\/span><span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'clusterTime'<\/span><span class=\"p\">].<\/span><span class=\"nf\">time<\/span><span class=\"p\">).<\/span><span class=\"nf\">isoformat<\/span><span class=\"p\">()<\/span>\r\n                        <span class=\"n\">payload<\/span> <span class=\"o\">=<\/span> <span class=\"p\">{<\/span><span class=\"s1\">'_id'<\/span><span class=\"ss\">:doc_id<\/span><span class=\"p\">}<\/span>\r\n                        <span class=\"c1\">######## Uncomment the following line if you want to add operation metadata fields to the document event. <\/span>\r\n                        <span class=\"n\">payload<\/span><span class=\"p\">.<\/span><span class=\"nf\">update<\/span><span class=\"p\">({<\/span><span class=\"s1\">'action'<\/span><span class=\"ss\">:op_type<\/span><span class=\"p\">,<\/span><span class=\"s1\">'timestamp'<\/span><span class=\"ss\">:str<\/span><span class=\"p\">(<\/span><span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'clusterTime'<\/span><span class=\"p\">].<\/span><span class=\"nf\">time<\/span><span class=\"p\">),<\/span><span class=\"s1\">'timestampReadable'<\/span><span class=\"ss\">:str<\/span><span class=\"p\">(<\/span><span class=\"n\">readable<\/span><span class=\"p\">)})<\/span>\r\n                        <span class=\"c1\">######## Uncomment the following line if you want to add db and coll metadata fields to the document event. <\/span>\r\n                        <span class=\"n\">payload<\/span><span class=\"p\">.<\/span><span class=\"nf\">update<\/span><span class=\"p\">({<\/span><span class=\"s1\">'db'<\/span><span class=\"ss\">:str<\/span><span class=\"p\">(<\/span><span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'ns'<\/span><span class=\"p\">][<\/span><span class=\"s1\">'db'<\/span><span class=\"p\">]),<\/span><span class=\"s1\">'coll'<\/span><span class=\"ss\">:str<\/span><span class=\"p\">(<\/span><span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'ns'<\/span><span class=\"p\">][<\/span><span class=\"s1\">'coll'<\/span><span class=\"p\">])})<\/span>\r\n\r\n                        <span class=\"c1\"># Append event for S3<\/span>\r\n                        <span class=\"k\">if<\/span> <span class=\"s2\">\"BUCKET_NAME\"<\/span> <span class=\"k\">in<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">:<\/span>\r\n                            <span class=\"n\">put_s3_event<\/span><span class=\"p\">(<\/span><span class=\"n\">json_util<\/span><span class=\"p\">.<\/span><span class=\"nf\">dumps<\/span><span class=\"p\">(<\/span><span class=\"n\">payload<\/span><span class=\"p\">),<\/span> <span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'ns'<\/span><span class=\"p\">][<\/span><span class=\"s1\">'db'<\/span><span class=\"p\">]),<\/span> <span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">change_event<\/span><span class=\"p\">[<\/span><span class=\"s1\">'ns'<\/span><span class=\"p\">][<\/span><span class=\"s1\">'coll'<\/span><span class=\"p\">]),<\/span><span class=\"n\">op_id<\/span><span class=\"p\">)<\/span>\r\n\r\n                        <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Processed event ID {}'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">op_id<\/span><span class=\"p\">))<\/span>\r\n\r\n                    <span class=\"n\">events_processed<\/span> <span class=\"o\">+=<\/span> <span class=\"mi\">1<\/span>\r\n\r\n                    <span class=\"k\">if<\/span> <span class=\"n\">events_processed<\/span> <span class=\"o\">&gt;=<\/span> <span class=\"n\">state_sync_count<\/span> <span class=\"ow\">and<\/span> <span class=\"s2\">\"BUCKET_NAME\"<\/span> <span class=\"ow\">not<\/span> <span class=\"k\">in<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">:<\/span>\r\n                        <span class=\"c1\"># To reduce DocumentDB IO, only persist the stream state every N events<\/span>\r\n                        <span class=\"n\">store_last_processed_id<\/span><span class=\"p\">(<\/span><span class=\"n\">change_stream<\/span><span class=\"p\">.<\/span><span class=\"nf\">resume_token<\/span><span class=\"p\">)<\/span>\r\n                        <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Synced token {} to state collection'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">change_stream<\/span><span class=\"p\">.<\/span><span class=\"nf\">resume_token<\/span><span class=\"p\">))<\/span>\r\n\r\n    <span class=\"n\">except<\/span> <span class=\"no\">OperationFailure<\/span> <span class=\"n\">as<\/span> <span class=\"ss\">of:\r\n        <\/span><span class=\"n\">send_sns_alert<\/span><span class=\"p\">(<\/span><span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">of<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"k\">if<\/span> <span class=\"n\">of<\/span><span class=\"p\">.<\/span><span class=\"nf\">code<\/span> <span class=\"o\">==<\/span> <span class=\"no\">TOKEN_DATA_DELETED_CODE<\/span><span class=\"p\">:<\/span>\r\n            <span class=\"c1\"># Data for the last processed ID has been deleted in the change stream,<\/span>\r\n            <span class=\"c1\"># Store the last known good state so our next invocation<\/span>\r\n            <span class=\"c1\"># starts from the most recently available data<\/span>\r\n            <span class=\"n\">store_last_processed_id<\/span><span class=\"p\">(<\/span><span class=\"no\">None<\/span><span class=\"p\">)<\/span>\r\n        <span class=\"k\">raise<\/span>\r\n\r\n    <span class=\"n\">except<\/span> <span class=\"no\">Exception<\/span> <span class=\"n\">as<\/span> <span class=\"ss\">ex:\r\n        <\/span><span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">error<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Exception in executing replication: {}'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"n\">send_sns_alert<\/span><span class=\"p\">(<\/span><span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">ex<\/span><span class=\"p\">))<\/span>\r\n        <span class=\"k\">raise<\/span>\r\n\r\n    <span class=\"ss\">else:\r\n        \r\n        <\/span><span class=\"k\">if<\/span> <span class=\"n\">events_processed<\/span> <span class=\"o\">&gt;<\/span> <span class=\"mi\">0<\/span><span class=\"p\">:<\/span>\r\n\r\n            <span class=\"n\">store_last_processed_id<\/span><span class=\"p\">(<\/span><span class=\"n\">change_stream<\/span><span class=\"p\">.<\/span><span class=\"nf\">resume_token<\/span><span class=\"p\">)<\/span>\r\n            <span class=\"n\">logger<\/span><span class=\"p\">.<\/span><span class=\"nf\">debug<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Synced token {} to state collection'<\/span><span class=\"p\">.<\/span><span class=\"nf\">format<\/span><span class=\"p\">(<\/span><span class=\"n\">change_stream<\/span><span class=\"p\">.<\/span><span class=\"nf\">resume_token<\/span><span class=\"p\">))<\/span>\r\n            <span class=\"k\">return<\/span><span class=\"p\">{<\/span>\r\n                <span class=\"s1\">'statusCode'<\/span><span class=\"p\">:<\/span> <span class=\"mi\">200<\/span><span class=\"p\">,<\/span>\r\n                <span class=\"s1\">'description'<\/span><span class=\"p\">:<\/span> <span class=\"s1\">'Success'<\/span><span class=\"p\">,<\/span>\r\n                <span class=\"s1\">'detail'<\/span><span class=\"p\">:<\/span> <span class=\"n\">json<\/span><span class=\"p\">.<\/span><span class=\"nf\">dumps<\/span><span class=\"p\">(<\/span><span class=\"n\">str<\/span><span class=\"p\">(<\/span><span class=\"n\">events_processed<\/span><span class=\"p\">)<\/span><span class=\"o\">+<\/span> <span class=\"s1\">' records processed successfully.'<\/span><span class=\"p\">)<\/span>\r\n            <span class=\"p\">}<\/span>\r\n        <span class=\"ss\">else:\r\n            <\/span><span class=\"k\">if<\/span> <span class=\"n\">canary_record<\/span> <span class=\"n\">is<\/span> <span class=\"ow\">not<\/span> <span class=\"no\">None<\/span><span class=\"p\">:<\/span>\r\n                <span class=\"k\">return<\/span><span class=\"p\">{<\/span>\r\n                    <span class=\"s1\">'statusCode'<\/span><span class=\"p\">:<\/span> <span class=\"mi\">202<\/span><span class=\"p\">,<\/span>\r\n                    <span class=\"s1\">'description'<\/span><span class=\"p\">:<\/span> <span class=\"s1\">'Success'<\/span><span class=\"p\">,<\/span>\r\n                    <span class=\"s1\">'detail'<\/span><span class=\"p\">:<\/span> <span class=\"n\">json<\/span><span class=\"p\">.<\/span><span class=\"nf\">dumps<\/span><span class=\"p\">(<\/span><span class=\"s1\">'Canary applied. No records to process.'<\/span><span class=\"p\">)<\/span>\r\n                <span class=\"p\">}<\/span>\r\n            <span class=\"ss\">else:\r\n                <\/span><span class=\"k\">return<\/span><span class=\"p\">{<\/span>\r\n                    <span class=\"s1\">'statusCode'<\/span><span class=\"p\">:<\/span> <span class=\"mi\">201<\/span><span class=\"p\">,<\/span>\r\n                    <span class=\"s1\">'description'<\/span><span class=\"p\">:<\/span> <span class=\"s1\">'Success'<\/span><span class=\"p\">,<\/span>\r\n                    <span class=\"s1\">'detail'<\/span><span class=\"p\">:<\/span> <span class=\"n\">json<\/span><span class=\"p\">.<\/span><span class=\"nf\">dumps<\/span><span class=\"p\">(<\/span><span class=\"s1\">'No records to process.'<\/span><span class=\"p\">)<\/span>\r\n                <span class=\"p\">}<\/span>\r\n\r\n    <span class=\"ss\">finally:\r\n\r\n        <\/span><span class=\"c1\"># Close Kafka client<\/span>\r\n        <span class=\"k\">if<\/span> <span class=\"s2\">\"MSK_BOOTSTRAP_SRV\"<\/span> <span class=\"k\">in<\/span> <span class=\"n\">os<\/span><span class=\"p\">.<\/span><span class=\"nf\">environ<\/span><span class=\"p\">:<\/span>                                                 \r\n            <span class=\"n\">kafka_client<\/span><span class=\"p\">.<\/span><span class=\"nf\">close<\/span><span class=\"p\">()<\/span>\r\n\r\n\r\n<\/code><\/pre>\n","protected":false},"excerpt":{"rendered":"<p>&#8212; #!\/bin\/env python import json import logging import [&hellip;]<\/p>\n","protected":false},"author":9,"featured_media":0,"comment_status":"closed","ping_status":"open","sticky":false,"template":"","format":"standard","meta":{"footnotes":""},"categories":[1],"tags":[],"class_list":["post-47087","post","type-post","status-publish","format-standard","hentry","category-uncategorized"],"yoast_head":"<!-- This site is optimized with the Yoast SEO Premium plugin v21.5 (Yoast SEO v21.5) - https:\/\/yoast.com\/wordpress\/plugins\/seo\/ -->\n<title>2022\u5e7410\u670819\u65e5\u7684CloudFormation CodePipeline\u5907\u5fd8\u5f55 - Blog - Silicon Cloud<\/title>\n<meta name=\"robots\" content=\"index, follow, max-snippet:-1, max-image-preview:large, max-video-preview:-1\" \/>\n<link rel=\"canonical\" href=\"https:\/\/www.silicloud.com\/zh\/blog\/2022\u5e7410\u670819\u65e5\u7684cloudformation-codepipeline\u5907\u5fd8\u5f55\u3002\/\" \/>\n<meta property=\"og:locale\" content=\"zh_CN\" \/>\n<meta property=\"og:type\" content=\"article\" \/>\n<meta property=\"og:title\" content=\"2022\u5e7410\u670819\u65e5\u7684CloudFormation CodePipeline\u5907\u5fd8\u5f55\" \/>\n<meta property=\"og:description\" content=\"--- #!\/bin\/env python import json import logging import [&hellip;]\" \/>\n<meta property=\"og:url\" content=\"https:\/\/www.silicloud.com\/zh\/blog\/2022\u5e7410\u670819\u65e5\u7684cloudformation-codepipeline\u5907\u5fd8\u5f55\u3002\/\" \/>\n<meta property=\"og:site_name\" content=\"Blog - Silicon Cloud\" \/>\n<meta property=\"article:published_time\" content=\"2022-12-29T02:13:02+00:00\" \/>\n<meta property=\"article:modified_time\" content=\"2024-01-15T04:16:44+00:00\" \/>\n<meta name=\"author\" content=\"\u6e05, \u626c\" \/>\n<meta name=\"twitter:card\" content=\"summary_large_image\" \/>\n<meta name=\"twitter:label1\" content=\"\u4f5c\u8005\" \/>\n\t<meta name=\"twitter:data1\" content=\"\u6e05, \u626c\" \/>\n\t<meta name=\"twitter:label2\" content=\"\u9884\u8ba1\u9605\u8bfb\u65f6\u95f4\" \/>\n\t<meta name=\"twitter:data2\" content=\"10 \u5206\" \/>\n<script type=\"application\/ld+json\" class=\"yoast-schema-graph\">{\"@context\":\"https:\/\/schema.org\",\"@graph\":[{\"@type\":\"WebPage\",\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/2022%e5%b9%b410%e6%9c%8819%e6%97%a5%e7%9a%84cloudformation-codepipeline%e5%a4%87%e5%bf%98%e5%bd%95%e3%80%82\/\",\"url\":\"https:\/\/www.silicloud.com\/zh\/blog\/2022%e5%b9%b410%e6%9c%8819%e6%97%a5%e7%9a%84cloudformation-codepipeline%e5%a4%87%e5%bf%98%e5%bd%95%e3%80%82\/\",\"name\":\"2022\u5e7410\u670819\u65e5\u7684CloudFormation CodePipeline\u5907\u5fd8\u5f55 - Blog - Silicon Cloud\",\"isPartOf\":{\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/#website\"},\"datePublished\":\"2022-12-29T02:13:02+00:00\",\"dateModified\":\"2024-01-15T04:16:44+00:00\",\"author\":{\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/#\/schema\/person\/cb5556d2501da73d864cac945e8d9461\"},\"breadcrumb\":{\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/2022%e5%b9%b410%e6%9c%8819%e6%97%a5%e7%9a%84cloudformation-codepipeline%e5%a4%87%e5%bf%98%e5%bd%95%e3%80%82\/#breadcrumb\"},\"inLanguage\":\"zh-Hans\",\"potentialAction\":[{\"@type\":\"ReadAction\",\"target\":[\"https:\/\/www.silicloud.com\/zh\/blog\/2022%e5%b9%b410%e6%9c%8819%e6%97%a5%e7%9a%84cloudformation-codepipeline%e5%a4%87%e5%bf%98%e5%bd%95%e3%80%82\/\"]}]},{\"@type\":\"BreadcrumbList\",\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/2022%e5%b9%b410%e6%9c%8819%e6%97%a5%e7%9a%84cloudformation-codepipeline%e5%a4%87%e5%bf%98%e5%bd%95%e3%80%82\/#breadcrumb\",\"itemListElement\":[{\"@type\":\"ListItem\",\"position\":1,\"name\":\"\u9996\u9875\",\"item\":\"https:\/\/www.silicloud.com\/zh\/blog\/\"},{\"@type\":\"ListItem\",\"position\":2,\"name\":\"2022\u5e7410\u670819\u65e5\u7684CloudFormation CodePipeline\u5907\u5fd8\u5f55\"}]},{\"@type\":\"WebSite\",\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/#website\",\"url\":\"https:\/\/www.silicloud.com\/zh\/blog\/\",\"name\":\"Blog - Silicon Cloud\",\"description\":\"\",\"inLanguage\":\"zh-Hans\"},{\"@type\":\"Person\",\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/#\/schema\/person\/cb5556d2501da73d864cac945e8d9461\",\"name\":\"\u6e05, \u626c\",\"image\":{\"@type\":\"ImageObject\",\"inLanguage\":\"zh-Hans\",\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/#\/schema\/person\/image\/\",\"url\":\"https:\/\/secure.gravatar.com\/avatar\/32a4239de8ff29adace466261d309424a1e5fe9f7e3036bf89fe03f2e3dbe717?s=96&d=mm&r=g\",\"contentUrl\":\"https:\/\/secure.gravatar.com\/avatar\/32a4239de8ff29adace466261d309424a1e5fe9f7e3036bf89fe03f2e3dbe717?s=96&d=mm&r=g\",\"caption\":\"\u6e05, \u626c\"},\"url\":\"https:\/\/www.silicloud.com\/zh\/blog\/author\/qingyang\/\"},{\"@type\":\"ImageObject\",\"inLanguage\":\"zh-Hans\",\"@id\":\"https:\/\/www.silicloud.com\/zh\/blog\/2022%e5%b9%b410%e6%9c%8819%e6%97%a5%e7%9a%84cloudformation-codepipeline%e5%a4%87%e5%bf%98%e5%bd%95%e3%80%82\/#local-main-organization-logo\",\"url\":\"\",\"contentUrl\":\"\",\"caption\":\"Blog - Silicon Cloud\"}]}<\/script>\n<!-- \/ Yoast SEO Premium plugin. -->","yoast_head_json":{"title":"2022\u5e7410\u670819\u65e5\u7684CloudFormation CodePipeline\u5907\u5fd8\u5f55 - Blog - Silicon Cloud","robots":{"index":"index","follow":"follow","max-snippet":"max-snippet:-1","max-image-preview":"max-image-preview:large","max-video-preview":"max-video-preview:-1"},"canonical":"https:\/\/www.silicloud.com\/zh\/blog\/2022\u5e7410\u670819\u65e5\u7684cloudformation-codepipeline\u5907\u5fd8\u5f55\u3002\/","og_locale":"zh_CN","og_type":"article","og_title":"2022\u5e7410\u670819\u65e5\u7684CloudFormation CodePipeline\u5907\u5fd8\u5f55","og_description":"--- #!\/bin\/env python import json import logging import [&hellip;]","og_url":"https:\/\/www.silicloud.com\/zh\/blog\/2022\u5e7410\u670819\u65e5\u7684cloudformation-codepipeline\u5907\u5fd8\u5f55\u3002\/","og_site_name":"Blog - Silicon Cloud","article_published_time":"2022-12-29T02:13:02+00:00","article_modified_time":"2024-01-15T04:16:44+00:00","author":"\u6e05, \u626c","twitter_card":"summary_large_image","twitter_misc":{"\u4f5c\u8005":"\u6e05, \u626c","\u9884\u8ba1\u9605\u8bfb\u65f6\u95f4":"10 \u5206"},"schema":{"@context":"https:\/\/schema.org","@graph":[{"@type":"WebPage","@id":"https:\/\/www.silicloud.com\/zh\/blog\/2022%e5%b9%b410%e6%9c%8819%e6%97%a5%e7%9a%84cloudformation-codepipeline%e5%a4%87%e5%bf%98%e5%bd%95%e3%80%82\/","url":"https:\/\/www.silicloud.com\/zh\/blog\/2022%e5%b9%b410%e6%9c%8819%e6%97%a5%e7%9a%84cloudformation-codepipeline%e5%a4%87%e5%bf%98%e5%bd%95%e3%80%82\/","name":"2022\u5e7410\u670819\u65e5\u7684CloudFormation CodePipeline\u5907\u5fd8\u5f55 - Blog - Silicon Cloud","isPartOf":{"@id":"https:\/\/www.silicloud.com\/zh\/blog\/#website"},"datePublished":"2022-12-29T02:13:02+00:00","dateModified":"2024-01-15T04:16:44+00:00","author":{"@id":"https:\/\/www.silicloud.com\/zh\/blog\/#\/schema\/person\/cb5556d2501da73d864cac945e8d9461"},"breadcrumb":{"@id":"https:\/\/www.silicloud.com\/zh\/blog\/2022%e5%b9%b410%e6%9c%8819%e6%97%a5%e7%9a%84cloudformation-codepipeline%e5%a4%87%e5%bf%98%e5%bd%95%e3%80%82\/#breadcrumb"},"inLanguage":"zh-Hans","potentialAction":[{"@type":"ReadAction","target":["https:\/\/www.silicloud.com\/zh\/blog\/2022%e5%b9%b410%e6%9c%8819%e6%97%a5%e7%9a%84cloudformation-codepipeline%e5%a4%87%e5%bf%98%e5%bd%95%e3%80%82\/"]}]},{"@type":"BreadcrumbList","@id":"https:\/\/www.silicloud.com\/zh\/blog\/2022%e5%b9%b410%e6%9c%8819%e6%97%a5%e7%9a%84cloudformation-codepipeline%e5%a4%87%e5%bf%98%e5%bd%95%e3%80%82\/#breadcrumb","itemListElement":[{"@type":"ListItem","position":1,"name":"\u9996\u9875","item":"https:\/\/www.silicloud.com\/zh\/blog\/"},{"@type":"ListItem","position":2,"name":"2022\u5e7410\u670819\u65e5\u7684CloudFormation CodePipeline\u5907\u5fd8\u5f55"}]},{"@type":"WebSite","@id":"https:\/\/www.silicloud.com\/zh\/blog\/#website","url":"https:\/\/www.silicloud.com\/zh\/blog\/","name":"Blog - Silicon Cloud","description":"","inLanguage":"zh-Hans"},{"@type":"Person","@id":"https:\/\/www.silicloud.com\/zh\/blog\/#\/schema\/person\/cb5556d2501da73d864cac945e8d9461","name":"\u6e05, \u626c","image":{"@type":"ImageObject","inLanguage":"zh-Hans","@id":"https:\/\/www.silicloud.com\/zh\/blog\/#\/schema\/person\/image\/","url":"https:\/\/secure.gravatar.com\/avatar\/32a4239de8ff29adace466261d309424a1e5fe9f7e3036bf89fe03f2e3dbe717?s=96&d=mm&r=g","contentUrl":"https:\/\/secure.gravatar.com\/avatar\/32a4239de8ff29adace466261d309424a1e5fe9f7e3036bf89fe03f2e3dbe717?s=96&d=mm&r=g","caption":"\u6e05, \u626c"},"url":"https:\/\/www.silicloud.com\/zh\/blog\/author\/qingyang\/"},{"@type":"ImageObject","inLanguage":"zh-Hans","@id":"https:\/\/www.silicloud.com\/zh\/blog\/2022%e5%b9%b410%e6%9c%8819%e6%97%a5%e7%9a%84cloudformation-codepipeline%e5%a4%87%e5%bf%98%e5%bd%95%e3%80%82\/#local-main-organization-logo","url":"","contentUrl":"","caption":"Blog - Silicon Cloud"}]}},"_links":{"self":[{"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/posts\/47087","targetHints":{"allow":["GET"]}}],"collection":[{"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/posts"}],"about":[{"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/types\/post"}],"author":[{"embeddable":true,"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/users\/9"}],"replies":[{"embeddable":true,"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/comments?post=47087"}],"version-history":[{"count":1,"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/posts\/47087\/revisions"}],"predecessor-version":[{"id":54914,"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/posts\/47087\/revisions\/54914"}],"wp:attachment":[{"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/media?parent=47087"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/categories?post=47087"},{"taxonomy":"post_tag","embeddable":true,"href":"https:\/\/www.silicloud.com\/zh\/blog\/wp-json\/wp\/v2\/tags?post=47087"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}