Search Options

Results per page
Sort
Preferred Languages
Advance

Results 1 - 2 of 2 for Blog (0.13 sec)

  1. internal/event/target/kafka.go

    	objectName, err := url.QueryUnescape(eventData.S3.Object.Key)
    	if err != nil {
    		return nil, err
    	}
    
    	key := eventData.S3.Bucket.Name + "/" + objectName
    	data, err := json.Marshal(event.Log{EventName: eventData.EventName, Key: key, Records: []event.Event{eventData}})
    	if err != nil {
    		return nil, err
    	}
    
    	return &sarama.ProducerMessage{
    		Topic: target.args.Topic,
    Go
    - Registered: Sun May 05 19:28:20 GMT 2024
    - Last Modified: Tue Feb 20 08:16:35 GMT 2024
    - 13K bytes
    - Viewed (0)
  2. internal/logger/target/kafka/kafka.go

    			return ctx.Err()
    		}
    		return nil
    	default:
    		// log channel is full, do not wait and return
    		// an error immediately to the caller
    		atomic.AddInt64(&h.totalMessages, 1)
    		atomic.AddInt64(&h.failedMessages, 1)
    		return errors.New("log buffer full")
    	}
    	return nil
    }
    
    // SendFromStore - reads the log from store and sends it to kafka.
    func (h *Target) SendFromStore(key store.Key) (err error) {
    Go
    - Registered: Sun May 05 19:28:20 GMT 2024
    - Last Modified: Thu Feb 22 06:26:06 GMT 2024
    - 10.1K bytes
    - Viewed (1)
Back to top