Hello, I want to implement 2 endpoints. ```/bills ...
# help-connector-development
l
Hello, I want to implement 2 endpoints.
Copy code
/bills
/bills/{bill_id}
The first one is easy. How do I implement the 2nd using Python CDK?
k
A message from kapa.ai
To implement the second endpoint
/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:
Copy code
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).
l
@kapa.ai Any connector that is using the above concept?
k
A message from kapa.ai
Yes, there are several connectors that use the concepts outlined in the provided extracts. For instance, the [Square](https://docs.airbyte.com/integrations/sources/square) and [Notion](https://docs.airbyte.com/integrations/sources/notion) source connectors are examples of connectors that support various sync modes such as Full Refresh - Overwrite, Full Refresh - Append, Incremental - Append, and Incremental - Deduped History. Additionally, the [HTTP Request connector](https://docs.airbyte.com/integrations/sources/http-request) is an example of a connector that allows you to connect to any HTTP API. However, it's important to note that this connector is graveyarded and will not be receiving any updates from the Airbyte team. Its functionalities have been replaced by the [Airbyte CDK](https://docs.airbyte.com/connector-development/cdk-python/), which allows you to create source connectors for any HTTP API. For developing your own connectors, Airbyte provides some Connector Development Kits (CDKs) to help you build connectors. The [Airbyte Dotnet CDK](https://github.com/mrhamburg/airbyte.cdk.dotnet) is a framework for rapidly developing production-grade Airbyte connectors. The CDK currently offers helpers specific for creating Airbyte source connectors for HTTP APIs and generic Dotnet sources.
l
@kapa.ai can you provide connector examples where I can find parameterized endpoints built through Python cdk
k
A message from kapa.ai
Sure, the provided extracts mention a few examples of connectors built using the Python CDK that interact with parameterized endpoints. One such example is a connector that interacts with the [Congress API’s bill by congress endpoint](https://api.congress.gov/#/bill/bill_list_by_congress). This is a parameterized endpoint where the
congress
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.
l
Copy code
class 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?
e
@Lenin Mishra Did you ever figure this out? I am trying to do something similar and am stuck.
l
@Emily Allegood You can easily do it with Low code CDK using the Parent stream feature.