File tree Expand file tree Collapse file tree 1 file changed +7
-3
lines changed Expand file tree Collapse file tree 1 file changed +7
-3
lines changed Original file line number Diff line number Diff line change @@ -59,7 +59,7 @@ def carrier_get(key):
59
59
60
60
61
61
def _dsm_set_kinesis_context (event ):
62
- from ddtrace .internal . datastreams . botocore import calculate_kinesis_payload_size
62
+ from ddtrace .data_streams import set_consume_checkpoint
63
63
64
64
records = event .get ("Records" )
65
65
if records is None :
@@ -68,9 +68,13 @@ def _dsm_set_kinesis_context(event):
68
68
for record in records :
69
69
arn = record .get ("eventSourceARN" , "" )
70
70
context_json = _get_dsm_context_from_lambda (record )
71
- payload_size = calculate_kinesis_payload_size (record , context_json )
71
+ if not context_json :
72
+ logger .debug ("DataStreams skipped lambda message: %r" , record )
73
+ return None
72
74
73
- _dsm_set_context_helper ("kinesis" , arn , payload_size , context_json )
75
+ def carrier_get (key ):
76
+ return context_json .get (key )
77
+ set_consume_checkpoint ("kinesis" , arn , carrier_get )
74
78
75
79
76
80
def _get_dsm_context_from_lambda (message ):
You can’t perform that action at this time.
0 commit comments