Lenin Mishra
06/30/2023, 6:37 AM/bills
/bills/{bill_id}
The first one is easy. How do I implement the 2nd using Python CDK?kapa.ai
06/30/2023, 6:37 AMkapa.ai
06/30/2023, 6:38 AM/bills/{bill_id} using the Python CDK, you would need to parameterize the endpoint. An example of such an endpoint is the Congress API’s bill by congress endpoint GET /bill/:congress. These endpoints are also valid synchronous HTTP endpoints. In the Python CDK, you would typically define this in your HttpStream class. The name of this stream will be consistent in the naming of the JSON schema and the HttpStream class. For example, if your endpoint was for a Bill, you might have bill.json and Bill respectively. Here's a simplified example: python class Bill(HttpStream): ... def path(self, stream_state=None, stream_slice=None, next_page_token=None): return "/bills/{bill_id}" In this case, you would need to replace {bill_id} with the actual ID of the bill you want to access. This could be done dynamically based on the data you're processing. Please note that this is a simplified example and you might need to adjust it based on your specific use case and the API you're working with. For more detailed information, you can refer to the [Python CDK documentation](https://docs.airbyte.com/connector-development/cdk-python) and the [Python CDK tutorial](https://docs.airbyte.com/connector-development/tutorials/cdk-tutorial-python-http/getting-started).Lenin Mishra
06/30/2023, 6:43 AMkapa.ai
06/30/2023, 6:43 AMkapa.ai
06/30/2023, 6:43 AMLenin Mishra
06/30/2023, 6:48 AMkapa.ai
06/30/2023, 6:48 AMkapa.ai
06/30/2023, 6:48 AMcongress parameter is part of the URL path. However, the actual code or detailed implementation of this connector is not provided in the extracts. For a more detailed guide on how to build a connector using the Python CDK, you can refer to the [Python CDK Speedrun: Creating a Source](https://docs.airbyte.com/connector-development/tutorials/cdk-speedrun) tutorial and the [tutorial on building a connector to extract data from the Webflow API](https://airbyte.com/tutorials/extract-data-from-the-webflow-api). Please note that if the integration requires query params or body parameters that are not key-value pairs, you should use the Python CDK, as the Connector Builder currently only supports key-value query params and request body parameters.Lenin Mishra
06/30/2023, 10:09 AMclass ZohoBooksStream(HttpStream, ABC):
_REGION_TO_API_URL = MappingProxyType(
{
"US": "<https://www.zohoapis.com/books/v3/>",
"AU": "<https://www.zohoapis.com.au/books/v3/>",
"EU": "<https://www.zohoapis.eu/books/v3/>",
"IN": "<https://www.zohoapis.in/books/v3/>",
"CN": "<https://www.zohoapis.com.cn/books/v3/>",
"JP": "<https://www.zohoapis.jp/books/v3/>",
}
)
url_base = None
def next_page_token(self, response: requests.Response) -> Optional[Mapping[str, Any]]:
try:
json_response = response.json()
has_more_page = json_response.get('page_context', None).get('has_more_page', None)
if has_more_page:
print(f"Moving onto page - {int(json_response.get('page_context').get('page')) + 1}")
return int(json_response.get('page_context').get('page')) + 1
else:
print("No more page")
return None
except Exception as e:
return None
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]:
page = 1
if next_page_token:
page = next_page_token
else:
page = 1
return {"organization_id":{self.config['organization_id']}, "page":page}
def parse_response(self, response: requests.Response, **kwargs) -> Iterable[Mapping]:
return {}
class Bills(ZohoBooksStream):
# TODO: Fill in the primary key. Required. This is usually a unique field in the stream, like an ID or a timestamp.
primary_key = "bill_id"
def __init__(self, config, **args):
super(Bills, self).__init__(**args)
self.url_base = self._REGION_TO_API_URL[config['location'].upper()]
self.config = config
def path(
self, stream_state: Mapping[str, Any] = None, stream_slice: Mapping[str, Any] = None, next_page_token: Mapping[str, Any] = None
) -> str:
"""
TODO: Override this method to define the path this stream corresponds to. E.g. if the url is <https://example-api.com/v1/customers> then this
should return "customers". Required.
"""
return "bills"
def parse_response(self, response: requests.Response, **kwargs) -> Iterable[Mapping]:
"""
TODO: Override this method to define how a response is parsed.
:return an iterable containing each record in the response
"""
try:
if response.json().get('bills') != None:
for record in response.json()['bills']:
yield record
else:
print("No records")
yield {}
except Exception as e:
print(e)
class Bill(ZohoBooksStream):
# TODO: Fill in the primary key. Required. This is usually a unique field in the stream, like an ID or a timestamp.
primary_key = "bill_id"
def __init__(self, bill_id, config, **args):
super(Bill, self).__init__(**args)
self.url_base = self._REGION_TO_API_URL[config['location'].upper()]
self.config = config
self.bill_id = bill_id
def path(
self, stream_state: Mapping[str, Any] = None, stream_slice: Mapping[str, Any] = None, next_page_token: Mapping[str, Any] = None
) -> str:
"""
TODO: Override this method to define the path this stream corresponds to. E.g. if the url is <https://example-api.com/v1/customers> then this
should return "customers". Required.
"""
return f"bills/{self.bill_id}"
def parse_response(self, response: requests.Response, **kwargs) -> Iterable[Mapping]:
"""
TODO: Override this method to define how a response is parsed.
:return an iterable containing each record in the response
"""
try:
if response.json().get('bill') != None:
yield response.json().get('bill')
else:
print("No records")
yield {}
except Exception as e:
print(e)
class SourceZohoBooks(AbstractSource):
@staticmethod
def _get_bill_ids(authenticator: str = None, config: str = None) -> Mapping[str, str]:
"""
Most of the Webflow APIs require the collection id, but the streams that we are generating use the collection name.
This function will return a dictionary containing collection_name: collection_id entries.
"""
bills = []
bills_stream = Bills(authenticator=authenticator, config=config)
bills_record = bills_stream.read_records(sync_mode="full_refresh")
# Loop over the list of records and create a dictionary with name as key, and _id as value
try:
for collection_obj in list(bills_record):
bills.append(collection_obj["bill_id"])
return bills
except Exception as e:
print(e)
def check_connection(self, logger, config) -> Tuple[bool, any]:
"""
TODO: Implement a connection check to validate that the user-provided config can be used to connect to the underlying API
See <https://github.com/airbytehq/airbyte/blob/master/airbyte-integrations/connectors/source-stripe/source_stripe/source.py#L232>
for an example.
:param config: the user-input config object conforming to the connector's spec.yaml
:param logger: logger object
:return Tuple[bool, any]: (True, None) if the input config can be used to connect to the API successfully, (False, error) otherwise.
"""
locations = ['IN', 'EU', 'US', 'JP', 'AU']
input_location = config['location']
if input_location not in locations:
return False, f"Input location {input_location} is invalid. Please input one of the following locations: {locations}"
else:
return True, None
def streams(self, config: Mapping[str, Any]) -> List[Stream]:
"""
define streams here
"""
_REGION_TO_ACCESS_URL = MappingProxyType(
{
"US": "<https://accounts.zoho.com/oauth/v2/token>",
"AU": "<https://accounts.zoho.com.au/oauth/v2/token>",
"EU": "<https://accounts.zoho.eu/oauth/v2/token>?",
"IN": "<https://accounts.zoho.in/oauth/v2/token>",
"CN": "<https://accounts.zoho.com.cn/oauth/v2/token>",
"JP": "<https://accounts.zoho.jp/oauth/v2/token>",
}
)
client_id = config.get('client_id', None)
client_secret = config.get('client_secret', None)
refresh_token = config.get('refresh_token', None)
location = config.get('location', None)
token_refresh_endpoint = _REGION_TO_ACCESS_URL[config["location"].upper()]
auth = ZohoOauth2Authenticator(token_refresh_endpoint, client_id, client_secret, refresh_token, location) # Oauth2Authenticator is also available if you need oauth support
stream_endpoint = [ Bills(authenticator=auth, config=config)]
bill_ids = self._get_bill_ids(authenticator=auth, config=config)
for id in bill_ids:
stream_endpoint.append(Bill(authenticator=auth, config=config, bill_id=id))
print(stream_endpoint)
return stream_endpoint
I have added these lines of code. But Airbyte is throwing error. What is the issue here?Emily Allegood
12/12/2023, 5:30 PMLenin Mishra
12/22/2023, 2:11 PM