Harini Rajagopal
07/31/2023, 9:51 PM2023-07-31 21:41:52 INFO i.a.w.g.DefaultReplicationWorker(prepStateForLaterSaving):563 - Source did not output any state messages
2023-07-31 21:41:52 WARN i.a.w.g.DefaultReplicationWorker(prepStateForLaterSaving):571 - State capture: No new state, falling back on input state: io.airbyte.config.State@7e33b6c7[state={}]
Code updating and fetching cursor value and updating state
def request_params(
self,
stream_state: Mapping[str, Any],
stream_slice: Mapping[str, Any] = None,
next_page_token: Mapping[str, Any] = None,
) -> MutableMapping[str, Any]:
start_datetime = stream_slice["created_at"]
end_datetime= self.end_datetime
created_at__range = f"['{start_datetime}', '{end_datetime}']"
return {"created_at__range":created_at__range}
def stream_slices(
self, sync_mode,
cursor_field: List[str] = None,
stream_state: Mapping[str, Any] = None
) -> Iterable[Optional[Mapping[str, any]]]:
start_datetime = datetime.strptime(stream_state["created_at"], self.date_format) if stream_state and "created_at" in stream_state else self.start_datetime
return self._chunk_date_range(start_datetime)
def _chunk_date_range(self, start_datetime: datetime) -> List[Mapping[str, any]]:
dates=[]
dates.append({"created_at": start_datetime.strftime(self.date_format)})
return dates
@property
def state(self) -> Mapping[str, Any]:
date_format = '%Y-%m-%d %H:%M:%S.%f%z'
if self._cursor_value:
return {self.cursor_field: self._cursor_value.strftime(self.date_format)}
else:
return {self.cursor_field: self.start_datetime.strftime(self.date_format)}
@state.setter
def state(self, value: Mapping[str, Any]):
self._cursor_value = datetime.strptime(value[self.cursor_field], self.date_format)
def read_records(self, *args, **kwargs) -> Iterable[Mapping[str, Any]]:
for record in super().read_records(*args, **kwargs):
if len(record["data"]) == 0:
return
if self._cursor_value:
latest_record_date = max(datetime.strptime(data[self.cursor_field],self.date_format) for data in record["data"])
self._cursor_value = max(self._cursor_value, latest_record_date)
yield recordkapa.ai
07/31/2023, 9:52 PMread_records method not being able to update the _cursor_value correctly. In your read_records method, you're updating the _cursor_value only if it already exists. If it doesn't exist, it won't be updated. This could be the reason why the state is not being updated. Here's the relevant part of your code: python def read_records(self, *args, **kwargs) -> Iterable[Mapping[str, Any]]: for record in super().read_records(*args, **kwargs): if len(record["data"]) == 0: return if self._cursor_value: latest_record_date = max(datetime.strptime(data[self.cursor_field],self.date_format) for data in record["data"]) self._cursor_value = max(self._cursor_value, latest_record_date) yield record You might want to update this part to ensure _cursor_value is updated even if it doesn't exist initially. Here's a suggestion: python def read_records(self, *args, **kwargs) -> Iterable[Mapping[str, Any]]: for record in super().read_records(*args, **kwargs): if len(record["data"]) == 0: return latest_record_date = max(datetime.strptime(data[self.cursor_field],self.date_format) for data in record["data"]) if self._cursor_value: self._cursor_value = max(self._cursor_value, latest_record_date) else: self._cursor_value = latest_record_date yield record This way, _cursor_value will be updated with the latest record date even if it doesn't exist initially. Please note that this suggestion is based on the information available in the extracts and the code you've shared. If the issue persists, you might want to debug your code to ensure the read_records method is being called and the _cursor_value is being updated correctly.Harini Rajagopal
07/31/2023, 10:08 PMkapa.ai
07/31/2023, 10:09 PMSELECT * FROM state; command. This might give you some insights into why the state and cursor value are not getting updated. [Source](https://airbyte.com/tutorials/incremental-data-synchronization) If none of these suggestions work, I would recommend reaching out to the Airbyte community on their [Slack channel](https://airbytehq.slack.com) or [forum](https://discuss.airbyte.io) for further assistance.