|
| 1 | +```python Python Ingest v2 |
| 2 | +import os |
| 3 | + |
| 4 | +from unstructured_ingest.v2.pipeline.pipeline import Pipeline |
| 5 | +from unstructured_ingest.v2.interfaces import ProcessorConfig |
| 6 | + |
| 7 | +from unstructured_ingest.v2.processes.connectors.kafka.cloud import ( |
| 8 | + CloudKafkaIndexerConfig, |
| 9 | + CloudKafkaDownloaderConfig, |
| 10 | + CloudKafkaConnectionConfig, |
| 11 | + CloudKafkaAccessConfig |
| 12 | +) |
| 13 | + |
| 14 | +from unstructured_ingest.v2.processes.partitioner import PartitionerConfig |
| 15 | +from unstructured_ingest.v2.processes.chunker import ChunkerConfig |
| 16 | +from unstructured_ingest.v2.processes.embedder import EmbedderConfig |
| 17 | +from unstructured_ingest.v2.processes.connectors.local import LocalUploaderConfig |
| 18 | + |
| 19 | +# Chunking and embedding are optional. |
| 20 | + |
| 21 | +if __name__ == "__main__": |
| 22 | + Pipeline.from_configs( |
| 23 | + context=ProcessorConfig(), |
| 24 | + indexer_config=CloudKafkaIndexerConfig( |
| 25 | + topic=os.getenv("KAFKA_TOPIC"), |
| 26 | + num_messages_to_consume=100, |
| 27 | + timeout=1 |
| 28 | + ), |
| 29 | + downloader_config=CloudKafkaDownloaderConfig(download_dir=os.getenv("LOCAL_FILE_DOWNLOAD_DIR")), |
| 30 | + source_connection_config=CloudKafkaConnectionConfig( |
| 31 | + access_config=CloudKafkaAccessConfig( |
| 32 | + kafka_api_key=os.getenv("KAFKA_API_KEY"), |
| 33 | + secret=os.getenv("KAFKA_SECRET") |
| 34 | + ), |
| 35 | + bootstrap_server=os.getenv("KAFKA_BOOTSTRAP_SERVER"), |
| 36 | + port=os.getenv("KAFKA_PORT") |
| 37 | + ), |
| 38 | + partitioner_config=PartitionerConfig( |
| 39 | + partition_by_api=True, |
| 40 | + api_key=os.getenv("UNSTRUCTURED_API_KEY"), |
| 41 | + partition_endpoint=os.getenv("UNSTRUCTURED_API_URL"), |
| 42 | + additional_partition_args={ |
| 43 | + "split_pdf_page": True, |
| 44 | + "split_pdf_allow_failed": True, |
| 45 | + "split_pdf_concurrency_level": 15 |
| 46 | + } |
| 47 | + ), |
| 48 | + chunker_config=ChunkerConfig(chunking_strategy="by_title"), |
| 49 | + embedder_config=EmbedderConfig(embedding_provider="huggingface"), |
| 50 | + uploader_config=LocalUploaderConfig(output_dir=os.getenv("LOCAL_FILE_OUTPUT_DIR")) |
| 51 | + ).run() |
| 52 | +``` |
0 commit comments